MasstransferExporter/MasstransferInfrastructure/Process/Service/ProcessCommunicator.cs

86 lines
2.3 KiB
C#
Raw Normal View History

2024-06-19 07:53:34 +00:00
using MasstransferCommon.Utils;
using MasstransferCommunicate.Process.Client;
using MasstransferCommunicate.Process.Model;
using Serilog;
namespace MasstransferCommunicate.Process.Service;
/// <summary>
/// 进程间通讯器
/// </summary>
public class ProcessCommunicator
{
private static readonly Dictionary<string, List<Action<string>>> Subscribers = new();
private static ProcessHelper? _helper;
/// <summary>
/// 连接到服务端
/// </summary>
public static async Task Connect()
{
_helper = await ProcessHelper.CreateClient("Masstransfer");
_helper.MessageReceived += HandleMessageReceived;
}
/// <summary>
/// 订阅某个主题
/// </summary>
/// <param name="topic"></param>
/// <param name="delegate"></param>
public static Task<bool> Subscribe<T>(string topic, Action<string> @delegate)
{
if (!Subscribers.ContainsKey(topic))
{
Subscribers.Add(topic, []);
}
Subscribers[topic].Add(@delegate);
return Task.FromResult(true);
}
/// <summary>
/// 发送消息
/// </summary>
/// <param name="topic"></param>
/// <param name="message"></param>
public static async Task<bool> Send<T>(string topic, T message)
{
await _helper?.SendMessageAsync(JsonUtil.ToJson(new Message<T>(topic, message)))!;
return true;
}
/// <summary>
/// 处理接收到的消息
/// </summary>
/// <returns></returns>
private static void HandleMessageReceived(string? message)
{
if (message == null) return;
Console.WriteLine($"收到来自服务端的消息: {message}");
var dictionary = JsonUtil.ToDictionary(message);
if (dictionary == null) return;
var topic = dictionary["Topic"] as string;
var data = dictionary["Data"];
if (!Subscribers.TryGetValue(topic, out var subscribers)) return;
foreach (var subscriber in subscribers)
{
try
{
// 通知订阅者
subscriber(JsonUtil.ToJson(data));
}
catch (Exception exception)
{
Log.Error(exception, "订阅主题 {Topic} 时发生错误", topic);
}
}
}
}