120 lines
3.5 KiB
C#
120 lines
3.5 KiB
C#
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
|
||
|
|
{
|
||
|
|
/// <summary>
|
||
|
|
/// 已连接的客户端
|
||
|
|
/// key-ip 一个客户端仅允许一个连接
|
||
|
|
/// </summary>
|
||
|
|
|
||
|
|
public static readonly ConcurrentDictionary<string, IChannelHandlerContext> DicClientChannels = new();
|
||
|
|
|
||
|
|
/// <summary>
|
||
|
|
/// 移除有问题通道
|
||
|
|
/// </summary>
|
||
|
|
/// <param name="channel"></param>
|
||
|
|
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}");
|
||
|
|
}
|
||
|
|
}
|
||
|
|
|
||
|
|
/// <summary>
|
||
|
|
/// 指定客户端发送消息
|
||
|
|
/// </summary>
|
||
|
|
/// <param name="ip"></param>
|
||
|
|
/// <param name="msg"></param>
|
||
|
|
/// <returns></returns>
|
||
|
|
public async static Task<bool> 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;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// <summary>
|
||
|
|
/// 指定客户端发送byte数组
|
||
|
|
/// </summary>
|
||
|
|
/// <param name="ip"></param>
|
||
|
|
/// <param name="msg"></param>
|
||
|
|
/// <returns></returns>
|
||
|
|
public async static Task<bool> 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;
|
||
|
|
}
|
||
|
|
|
||
|
|
/// <summary>
|
||
|
|
/// 发送广播消息
|
||
|
|
/// </summary>
|
||
|
|
/// <param name="messagePackage"></param>
|
||
|
|
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}");
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}
|