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}"); } } } }