using DotNetty.Transport.Channels;
using JSMachine.WMS.Infrastructure.Helper;
using System;
using System.Collections.Concurrent;
using System.Collections.Generic;
using System.Data;
using System.Linq;
using System.Reactive.Linq;
using System.Threading.Tasks;
namespace JSMachine.WMS.Netty.NettyAsClient.Transport
{
public static class TransportManager
{
///
/// 已连接的客户端
/// key-ip 一个客户端仅允许一个连接
///
public static readonly ConcurrentDictionary DicClientChannels = new();
///
/// 移除有问题通道
///
///
public static void RemoveInActiveChannel(IChannelHandlerContext channel)
{
try
{
DicClientChannels
.Where(p => p.Value == channel)
.Select(p => p.Key)
.ToList()
?.ForEach(p => DicClientChannels.Remove(p, out IChannelHandlerContext channelHandlerContext));
}
catch (Exception ex)
{
LogHelper.Error($"移除有问题通道出错: {ex.Message} \n {ex.StackTrace}");
}
}
///
/// 指定客户端发送消息
///
///
///
///
public async static Task SendMsg(string ip, string msg)
{
try
{
if (!DicClientChannels.ContainsKey(ip))
{
LogHelper.Error($"{ip}连接失败");
return false;
}
IChannelHandlerContext channel = DicClientChannels[ip];
await channel.Channel.WriteAndFlushAsync(msg);
return true;
}
catch (Exception ex)
{
LogHelper.Error($"发送消息出错: {ex.Message} \n {ex.StackTrace}");
}
return false;
}
///
/// 指定客户端发送byte数组
///
///
///
///
public async static Task SendMsgByte(string ip, byte[] msg)
{
try
{
if (!DicClientChannels.ContainsKey(ip))
{
LogHelper.Error($"{ip}连接失败");
return false;
}
IChannelHandlerContext channel = DicClientChannels[ip];
await channel.Channel.WriteAndFlushAsync(msg);
return true;
}
catch (Exception ex)
{
LogHelper.Error($"发送消息出错: {ex.Message} \n {ex.StackTrace}");
}
return false;
}
///
/// 发送广播消息
///
///
public static void BroadcastMsg(string brodCastMsg)
{
try
{
DicClientChannels.Keys.ToList().ForEach(async key =>
{
await DicClientChannels[key].Channel.WriteAndFlushAsync(brodCastMsg);
});
}
catch (Exception ex)
{
LogHelper.Error($"发送广播消息出错: {ex.Message} \n {ex.StackTrace}");
}
}
}
}