using System.Collections.Concurrent;
using System.Net;
using System.Net.Sockets;
using System.Text;
using NewLife.Collections;
using NewLife.Data;
using NewLife.Log;
using NewLife.Messaging;
using NewLife.Model;
namespace NewLife.Net;
/// <summary>增强的UDP服务器</summary>
/// <remarks>
/// <para>封装了UDP服务器功能,支持同时作为服务端和客户端使用。</para>
/// <para>功能特性:</para>
/// <list type="bullet">
/// <item>支持UDP广播和组播</item>
/// <item>自动管理UDP会话</item>
/// <item>支持地址重用(快速重启)</item>
/// <item>支持环回数据过滤</item>
/// </list>
/// <para>接收模式二选一:默认自动接收(打开后启动接收环走事件推送,不允许拉取);打开前设置 AutoReceive=false 则只允许同步/异步拉取(直读Socket)。</para>
/// </remarks>
public class UdpServer : SessionBase, ISocketServer, ILogFeature
{
#region 属性
/// <summary>会话超时时间(秒)</summary>
/// <remarks>
/// 对于每一个会话连接,如果超过该时间仍然没有收到任何数据,则断开会话连接。
/// </remarks>
public Int32 SessionTimeout { get; set; }
/// <summary>是否接收来自自己广播的环回数据</summary>
/// <remarks>默认false,不接收自己发出的广播数据</remarks>
public Boolean Loopback { get; set; }
/// <summary>地址重用</summary>
/// <remarks>
/// <para>主要应用于网络服务器重启交替。默认false。</para>
/// <para>一个端口释放后会等待两分钟之后才能再被使用,SO_REUSEADDR是让端口释放后立即就可以被再次使用。</para>
/// <para>SO_REUSEADDR用于对TCP套接字处于TIME_WAIT状态下的socket,才可以重复绑定使用。</para>
/// </remarks>
public Boolean ReuseAddress { get; set; }
/// <summary>最大并行接收数。接收环并发待收的报文数量,默认CPU*1.6</summary>
/// <remarks>小于等于0时按1处理。仅控制接收环并行度,接收模式由 <see cref="SessionBase.AutoReceive"/> 决定</remarks>
public Int32 MaxAsync
{
get => MaxReceiveCount;
set => MaxReceiveCount = value > 0 ? value : 1;
}
#endregion
#region 构造
/// <summary>实例化增强UDP服务器</summary>
public UdpServer()
{
Local.Type = NetType.Udp;
Remote.Type = NetType.Udp;
_Sessions = new SessionCollection(this);
SessionTimeout = SocketSetting.Current.SessionTimeout;
// 处理UDP最大并发接收
MaxAsync = Environment.ProcessorCount * 16 / 10;
if (SocketSetting.Current.Debug) Log = XTrace.Log;
}
/// <summary>使用监听端口初始化</summary>
/// <param name="listenPort">监听端口</param>
public UdpServer(Int32 listenPort) : this() => Port = listenPort;
/// <summary>销毁资源</summary>
/// <param name="disposing">是否释放托管资源</param>
protected override void Dispose(Boolean disposing)
{
base.Dispose(disposing);
if (Active) Close(GetType().Name + (disposing ? "Dispose" : "GC"));
_Sessions.Dispose();
}
#endregion
#region 方法
/// <summary>打开</summary>
/// <param name="cancellationToken">取消通知</param>
protected override Task<Boolean> OnOpenAsync(CancellationToken cancellationToken)
{
var sock = Client;
if (sock == null || !sock.IsBound)
{
var uri = Remote;
// 根据目标地址适配本地IPv4/IPv6
if (Local.Address.IsAny() && uri != null && !uri.Address.IsAny())
{
Local.Address = Local.Address.GetRightAny(uri.Address.AddressFamily);
}
Client = sock = NetHelper.CreateUdp(Local.Address.IsIPv4());
try
{
// 地址重用,主要应用于网络服务器重启交替。前一个进程关闭时,端口在短时间内处于TIME_WAIT,导致新进程无法监听。
// 启用地址重用后,即使旧进程未退出,新进程也可以监听,但只有旧进程退出后,新进程才能接受对该端口的连接请求
if (ReuseAddress) sock.SetSocketOption(SocketOptionLevel.Socket, SocketOptionName.ReuseAddress, true);
// 无需设置SocketOptionName.PacketInformation:接收环已改用 ReceiveFromAsync,不依赖包信息
//if (sock.AddressFamily == AddressFamily.InterNetwork)
// sock.SetSocketOption(SocketOptionLevel.IP, SocketOptionName.PacketInformation, true);
}
catch (Exception ex)
{
// 有些平台不支持地址重用,比如旧版A2上的Ubuntu16,但是云服务器的Linux都支持
XTrace.WriteLine(ex.Message);
}
// 打开前本地端口为 0,说明未指定监听端口(纯客户端模式)。首次判定后保持,后续不再依赖已被赋值的端口
if (Local.Port == 0) _autoLocalPort = true;
sock.Bind(Local.EndPoint);
if (Local.Port == 0 && sock.LocalEndPoint is IPEndPoint ep)
Local.Port = ep.Port;
// 客户端模式连接远端:发送走已连接路径(免每包 SendTo 序列化 SocketAddress),接收也只收该远端数据。
// 指定了本地端口的实例可能同时监听(双角色),保持 SendTo 与全量接收语义
if (_autoLocalPort && !Runtime.Mono && Remote != null && Remote.EndPoint is IPEndPoint remoteEP)
TryConnect(sock, remoteEP);
WriteLog("Open {0}", this);
}
return Task.FromResult(true);
}
/// <summary>客户端模式连接远端,收发都限定到该端点。广播/组播地址不能连接,保持 SendTo 语义</summary>
/// <param name="sock">底层Socket</param>
/// <param name="remote">远端端点</param>
private void TryConnect(Socket sock, IPEndPoint remote)
{
if (remote.Port <= 0 || remote.Address.IsAny()) return;
var buf = remote.Address.GetAddressBytes();
if (buf.Length == 4)
{
// 广播与组播地址不能连接,否则无法再向该地址发送
if (buf[0] >= 224 && buf[0] <= 239) return;
if (buf[3] == 255) return;
}
else if (remote.Address.IsIPv6Multicast) return;
try
{
sock.Connect(remote);
}
catch (Exception ex)
{
// 少数平台不支持连接(如未开广播的子网地址),回退 SendTo 发送路径即可
WriteLog("Connect {0} 失败:{1}", remote, ex.Message);
}
}
/// <summary>关闭</summary>
/// <param name="reason">关闭原因。便于日志分析</param>
/// <param name="cancellationToken">取消通知</param>
protected override Task<Boolean> OnCloseAsync(String reason, CancellationToken cancellationToken)
{
var sock = Client;
if (sock != null)
{
WriteLog("Close {0} {1}", reason, this);
try
{
// 以客户端模式工作时,发空包通知服务端结束会话。
// 注意:本方法处于 Dispose 流程内(Disposed 已置位),不能用 Send——其守卫会直接抛 ObjectDisposedException,
// 而异常被本 catch 吞掉后,下面的收尾步骤全部跳过,套接字与本地端口泄漏。这里改走不带守卫的 OnSend。
var remote = Remote;
if (remote != null && !remote.Address.IsAny() && remote.Port != 0)
{
OnSend(Pool.Empty);
}
}
catch (Exception ex)
{
if (!ex.IsDisposed()) OnError("Close", ex);
//if (ThrowException) throw;
}
finally
{
// 收尾与发送成败无关:无论通知服务端是否成功,都必须释放套接字并清理会话
Client = null;
CloseAllSession();
sock.Shutdown();
}
}
return Task.FromResult(true);
}
#endregion
#region 发送
/// <summary>发送数据</summary>
/// <remarks>
/// 目标地址由<seealso cref="SessionBase.Remote"/>决定
/// </remarks>
/// <param name="data">数据包</param>
/// <returns>是否成功</returns>
protected override Int32 OnSend(IPacket data) => OnSend(data, Remote.EndPoint);
/// <summary>发送数据</summary>
/// <remarks>
/// 目标地址由<seealso cref="SessionBase.Remote"/>决定
/// </remarks>
/// <param name="data">数据包</param>
/// <returns>是否成功</returns>
protected override Int32 OnSend(ArraySegment<Byte> data) => OnSend(data, Remote.EndPoint);
/// <summary>发送数据</summary>
/// <remarks>
/// 目标地址由<seealso cref="SessionBase.Remote"/>决定
/// </remarks>
/// <param name="data">数据包</param>
/// <returns>是否成功</returns>
protected override Int32 OnSend(ReadOnlySpan<Byte> data) => OnSend(data, Remote.EndPoint);
internal Int32 OnSend(IPacket pk, IPEndPoint remote)
{
var count = pk.Total;
using var span = Tracer?.NewSpan($"net:{Name}:Send", count + "", count);
try
{
var rs = 0;
// 服务端关闭时 OnCloseAsync 先将 Client 置 null 再通知会话,存在短暂竞态窗口:
// 若此时仍有 IOCP 线程执行 OnReceive → Send,会触发 NullReferenceException 或 throw。
// 静默返回 -1 代替 throw,避免关闭期间产生无意义的错误日志。
if (Client is not { } sock) return -1;
lock (_sendLock)
{
// Linux+Mono 的Connected总是true,需要特殊处理
var connected = sock.Connected;
if (Runtime.Mono && Runtime.Linux)
{
try
{
var r = sock.RemoteEndPoint;
connected = r is IPEndPoint ep && !ep.Address.IsAny() && ep.Port > 0;
}
catch
{
connected = false;
}
}
if (connected && !sock.EnableBroadcast)
{
if (Log.Enable && LogSend) WriteLog("Send [{0}]: {1}", count, pk.ToHex(LogDataLength));
if (count == 0)
rs = sock.Send(Pool.Empty);
else if (pk.Next == null && pk.TryGetArray(out var segment))
rs = sock.Send(segment.Array!, segment.Offset, segment.Count, SocketFlags.None);
#if NETCOREAPP || NETSTANDARD2_1
else if (pk.TryGetSpan(out var data))
rs = sock.Send(data);
#endif
else
rs = sock.Send(pk.ToSegments(), SocketFlags.None);
}
else
{
sock.CheckBroadcast(remote.Address);
if (Log.Enable && LogSend) WriteLog("Send {2} [{0}]: {1}", count, pk.ToHex(LogDataLength), remote);
if (count == 0)
rs = sock.SendTo(Pool.Empty, remote);
else if (pk.Next == null && pk.TryGetArray(out var segment))
rs = sock.SendTo(segment.Array!, segment.Offset, segment.Count, SocketFlags.None, remote);
#if NET6_0_OR_GREATER
else if (pk.TryGetSpan(out var data))
rs = sock.SendTo(data, remote);
#endif
else
rs = sock.SendTo(pk.ToArray(), 0, count, SocketFlags.None, remote);
}
}
return rs;
}
catch (Exception ex)
{
// 发生异常时,全量数据写入埋点
span?.SetError(ex, pk);
if (!ex.IsDisposed())
{
OnError("Send", ex);
}
return -1;
}
}
internal Int32 OnSend(Byte[] data, Int32 offset, Int32 count, IPEndPoint remote)
{
#if NET6_0_OR_GREATER
return OnSend(new ReadOnlySpan<Byte>(data, offset, count), remote);
#else
return OnSend(new ArraySegment<Byte>(data, offset, count), remote);
#endif
}
internal Int32 OnSend(ArraySegment<Byte> data, IPEndPoint remote)
{
var count = data.Count;
var logCount = count > LogDataLength ? count : LogDataLength;
using var span = Tracer?.NewSpan($"net:{Name}:Send", count + "", count);
try
{
var rs = 0;
// 服务端关闭时 Client 先被置 null,静默返回 -1 代替 throw
if (Client is not { } sock) return -1;
lock (_sendLock)
{
if (sock.Connected && !sock.EnableBroadcast)
{
if (Log.Enable && LogSend) WriteLog("Send [{0}]: {1}", count, data.Array.ToHex(data.Offset, logCount));
rs = sock.Send(data.Array!, data.Offset, data.Count, SocketFlags.None);
}
else
{
sock.CheckBroadcast(remote.Address);
if (Log.Enable && LogSend) WriteLog("Send {2} [{0}]: {1}", count, data.Array.ToHex(data.Offset, logCount), remote);
rs = sock.SendTo(data.Array!, data.Offset, data.Count, SocketFlags.None, remote);
}
}
return rs;
}
catch (Exception ex)
{
// 发生异常时,全量数据写入埋点
span?.SetError(ex, data.Array.ToHex(data.Offset, data.Count));
if (!ex.IsDisposed())
{
OnError("Send", ex);
}
return -1;
}
}
internal Int32 OnSend(ReadOnlySpan<Byte> data, IPEndPoint remote)
{
var count = data.Length;
using var span = Tracer?.NewSpan($"net:{Name}:Send", count + "", count);
try
{
var rs = 0;
// 服务端关闭时 Client 先被置 null,静默返回 -1 代替 throw
if (Client is not { } sock) return -1;
lock (_sendLock)
{
if (sock.Connected && !sock.EnableBroadcast)
{
if (Log.Enable && LogSend) WriteLog("Send [{0}]: {1}", count, data.ToHex(LogDataLength));
#if NETCOREAPP || NETSTANDARD2_1_OR_GREATER
rs = sock.Send(data, SocketFlags.None);
#else
rs = sock.Send(data.ToArray(), SocketFlags.None);
#endif
}
else
{
sock.CheckBroadcast(remote.Address);
if (Log.Enable && LogSend) WriteLog("Send {2} [{0}]: {1}", count, data.ToHex(LogDataLength), remote);
#if NET6_0_OR_GREATER
rs = sock.SendTo(data, SocketFlags.None, remote);
#else
rs = sock.SendTo(data.ToArray(), SocketFlags.None, remote);
#endif
}
}
return rs;
}
catch (Exception ex)
{
// 发生异常时,全量数据写入埋点
span?.SetError(ex, data.ToHex());
if (!ex.IsDisposed())
{
OnError("Send", ex);
}
return -1;
}
}
/// <summary>发送消息并等待匹配的响应。按远端查找/创建会话,由会话完成收发与配对</summary>
/// <param name="message">消息</param>
/// <param name="cancellationToken">取消通知</param>
/// <returns>响应消息</returns>
/// <remarks>必须调用会话的发送,否则配对会失败:响应数据报按来源端点分派给会话,匹配队列在会话上</remarks>
public override ValueTask<Object> SendMessageAsync(Object message, CancellationToken cancellationToken = default) => CreateSession(null, Remote.EndPoint).SendMessageAsync(message, cancellationToken);
/// <summary>发送消息并等待匹配的响应。按远端查找/创建会话,由会话完成收发与配对</summary>
/// <param name="message">消息</param>
/// <param name="cancellationToken">取消通知</param>
/// <returns>响应消息</returns>
public override ValueTask<IMessage> SendMessageAsync(IMessage message, CancellationToken cancellationToken = default)
{
if (CreateSession(null, Remote.EndPoint) is not UdpSession session) throw new NotSupportedException($"会话类型不支持数据报请求-响应等待 [{Remote.EndPoint}]");
return session.SendMessageAsync(message, cancellationToken);
}
#endregion
#region 接收
internal override Boolean OnReceiveAsync(SocketAsyncEventArgs se)
{
// 不能发起接收(关闭窗口内 Client 已清、或已停止)时与 TcpSession 同口径抛 ODE,由基类释放本槽。
// 不能返回 false——false 的语义是“本次同步完成、se 已持有新状态”,会把上一次完成的陈旧
// 事件参数重新投递(关闭窗口内成环:同一报文被反复派发,甚至触发空数据关闭分支)
if (!Active || Client is not { } sock) throw new ObjectDisposedException(GetType().Name);
// 每次接收以后,这个会被设置为远程地址,这里重置一下,以防万一。
// 注意:运行时在完成回调中“就地”更新该端点对象,且对象会随接收结果成为会话的远程端点被长期保留,
// 因此每次重投都必须新建对象;不能复用同一实例重置,否则会改到其它会话的端点。
se.RemoteEndPoint = new IPEndPoint(IPAddress.Any.GetRightAny(Local.EndPoint.AddressFamily), 0);
// 在StarAgent中,此时可能收到广播包,SocketFlags是Broadcast,需要清空,否则报错“参考的对象类型不支持尝试的操作”
se.SocketFlags = SocketFlags.None;
//return Client.ReceiveFromAsync(se);
// 历史写法:Android/旧 Mono 不支持 ReceiveMessageFromAsync,只能走 ReceiveFromAsync。
// 现统一走 ReceiveFromAsync:ReceiveMessageFromAsync 的包信息(本地地址)在本库仅用于填 e.Local,
// 却让运行时每次完成回调多构造一个端点对象——实测每包 216→104 B 分配、包率 +14%、每包 CPU -1.3 µs。
// 本地地址改由 GetReceiveLocalAddress(取 Local.Address,与消息路径 e.Local = Local.Address 同口径)提供。
return sock.ReceiveFromAsync(se);
}
/// <summary>本轮数据的本地地址。接收环走 ReceiveFromAsync,无包信息,取服务器绑定地址</summary>
/// <remarks>与消息路径 <c>e.Local = Local.Address</c> 同口径。多宿服务器绑 Any 时不再逐包给出实际到包地址</remarks>
/// <param name="se">接收事件参数</param>
/// <returns>本地地址</returns>
protected override IPAddress? GetReceiveLocalAddress(SocketAsyncEventArgs se) => Local.Address;
/// <summary>预处理</summary>
/// <param name="pk">数据包</param>
/// <param name="local">接收数据的本地地址</param>
/// <param name="remote">远程地址</param>
/// <returns>将要处理该数据包的会话</returns>
internal protected override ISocketSession? OnPreReceive(IPacket pk, IPAddress local, IPEndPoint remote)
{
// 过滤自己广播的环回数据。放在这里,兼容UdpSession
if (!Loopback && remote.Port == Port)
{
if (!Local.Address.IsAny())
{
if (remote.Address.Equals(Local.Address)) return null;
}
else
{
foreach (var item in NetHelper.GetIPsWithCache())
{
if (remote.Address.Equals(item)) return null;
}
}
}
// 为该连接单独创建一个会话,方便直接通信
return CreateSession(local, remote);
}
/// <summary>处理收到的数据</summary>
/// <param name="e">接收事件参数</param>
protected override Boolean OnReceive(ReceivedEventArgs e)
{
var pk = e.Packet;
var remote = e.Remote;
// 为该连接单独创建一个会话,方便直接通信
var session = remote == null ? null : CreateSession(e.Local, remote);
// 数据直接转交给会话,不再经过事件,那样在会话较多时极为浪费资源
if (session is UdpSession us)
{
// 协议模式:消息在会话内分发并升格到服务器层(RaiseReceiveInternal),原始数据报不再重复抛出
var processed = us.OnReceive(e);
if (!processed) RaiseReceive(session, e);
}
else
{
// 没有匹配到任何会话时,才在这里显示日志。理论上不存在这个可能性
if (Log.Enable && LogReceive && pk != null) WriteLog("Recv [{0}]: {1}", pk.Length, pk.ToHex(LogDataLength));
if (session != null) RaiseReceive(session, e);
}
return true;
}
/// <summary>收到异常时如何处理。Udp服务端不能关闭服务器,仅关闭出问题的会话</summary>
/// <param name="se"></param>
/// <returns>是否当作异常处理并结束会话。恒为 false:让接收环重投本槽继续收</returns>
internal override Boolean OnReceiveError(SocketAsyncEventArgs se)
{
// 缓冲区不足时,加大
if (se.SocketError == SocketError.MessageSize && BufferSize < 1024 * 1024) BufferSize *= 2;
// 关闭出问题的会话(Reset/Aborted 一般来自 ICMP 端口不可达),服务器本身不关闭
if (se.SocketError is SocketError.ConnectionReset or SocketError.ConnectionAborted)
{
var sessions = _Sessions;
if (sessions != null)
{
var ep = se.RemoteEndPoint as IPEndPoint;
var ss = ep != null ? sessions.Get(ep) : null;
ss?.Dispose();
}
}
// 一律返回 false,让接收环重投本槽继续收。返回 true 会让 ProcessEvent 释放本槽,
// 而接收槽只在启动时按 MaxAsync 创建一次、销毁后不重建,累计耗尽(默认 CPU*1.6 个)后
// 整台服务器会静默停收所有数据报。单个数据报出错(含 MessageSize)不代表远端不可达
return false;
}
#endregion
#region 会话
/// <summary>新会话时触发</summary>
public event EventHandler<SessionEventArgs>? NewSession;
private readonly SessionCollection _Sessions;
/// <summary>会话集合。用地址端口作为标识,业务应用自己维持地址端口与业务主键的对应关系。</summary>
/// <remarks>
/// <para>会话在 Socket 层是否存在以本集合为准;出集合只有三条路:会话释放(OnDisposed 级联)、
/// 超时清理不活动会话、停机清空,调用方不要手工 Remove。</para>
/// </remarks>
public IDictionary<String, ISocketSession> Sessions => _Sessions;
// 广播会话按端口索引。并发字典:无锁快路径读取与加锁写入并存,普通字典会在并发读写时损坏结构
// 派生索引:与 _Sessions 同源同销(同一个 OnDisposed 回调里移除),会话是否存在仍以 _Sessions 为准
private readonly ConcurrentDictionary<Int32, ISocketSession> _broadcasts = [];
/// <summary>发送锁。同一监听Socket上的所有发送(含各 UdpSession)均经此串行化;不用 Socket 实例作锁,避免外部代码对同一对象加锁导致死锁</summary>
private readonly Object _sendLock = new();
/// <summary>惰性打开锁。首个数据报可能来自不同端点、由多个 IOCP 线程并发进入,用于把打开收敛为一次</summary>
private readonly Object _openLock = new();
/// <summary>本地端口由系统自动分配(未指定监听端口,纯客户端模式)。首次打开时判定,决定是否连接远端</summary>
private Boolean _autoLocalPort;
Int32 g_ID = 0;
/// <summary>创建会话</summary>
/// <param name="local">接收数据的本地地址</param>
/// <param name="remoteEP">远程地址</param>
/// <returns></returns>
public virtual ISocketSession CreateSession(IPAddress? local, IPEndPoint remoteEP)
{
if (Disposed) throw new ObjectDisposedException(GetType().Name);
var sessions = _Sessions;
//if (sessions == null) return null;
// 平均执行耗时260.80ns,其中55%花在sessions.Get上面,Get里面有加锁操作
if (!Active)
{
// 首包可能来自不同端点、由多个 IOCP 线程并发进入:入口检查与置位之间没有原子性
// (见 SessionBase.OpenAsync 备注),两个线程同时 Open 会各建一个 Socket 并 Bind 同一端口,
// 后者抛 AddressAlreadyInUse 而丢掉该数据报。双检锁收敛为只打开一次;已建会话的快路径不取锁
lock (_openLock)
{
if (!Active)
{
// 根据目标地址适配本地IPv4/IPv6
Local.Address = Local.Address.GetRightAny(remoteEP.AddressFamily);
if (!Open()) throw new InvalidOperationException($"Open {Local} error");
}
}
}
// 需要查找已有会话,已有会话不存在时才创建新会话
// 端点键快路径:免去每包拼接字符串键(IPEndPoint.ToString 的多次字符串分配)
var session = sessions.Get(remoteEP);
// 是否匹配广播端口
var port = remoteEP.Port;
if (session != null || _broadcasts.TryGetValue(port, out session)) return session;
// 相同远程地址可能同时发来多个数据包,而底层采取多线程方式同时调度,导致创建多个会话
lock (sessions)
{
// 需要查找已有会话,已有会话不存在时才创建新会话
session = sessions.Get(remoteEP);
if (session != null || _broadcasts.TryGetValue(port, out session)) return session;
var us = new UdpSession(this, local, remoteEP)
{
Log = Log,
LogSend = LogSend,
LogReceive = LogReceive,
Tracer = Tracer,
};
session = us;
//us.ID = g_ID++;
// 会话改为原子操作,避免多线程冲突
us.ID = Interlocked.Increment(ref g_ID);
us.Tracer = Tracer;
try
{
// 必须在加入会话集合前完成启动和事件订阅。
// sessions.Add 会将会话插入 ConcurrentDictionary,随后其他 IOCP 线程可在 lock 外通过 sessions.Get 找到该会话并立即调用 OnReceive。
// 若此时 us.Start() 尚未执行(Received 事件未订阅),则首批数据包会静默丢弃。
us.Start();
// 触发新会话事件(用户代码如 EchoSession 在此处通过 NewSession 订阅 Ss_Received)
NewSession?.Invoke(this, new SessionEventArgs(session));
}
catch
{
// 从创建到入集合之间,会话还没被集合接管(Get 找不到它),启动或事件订阅抛异常时,
// 再也没人会来释放它:会话本身泄漏,同一端点的每个数据报还会再新建一个,持续累积。
// 此处立即释放后上抛,由接收环按既有错误处理记日志并继续收包(UdpSession 只解除与服务器的关联,不关共享Socket)
us.TryDispose();
throw;
}
if (sessions.Add(session))
{
// 广播地址,接受任何地址响应数据
if (Equals(remoteEP.Address, IPAddress.Broadcast))
{
_broadcasts[port] = session;
session.OnDisposed += (s, e) =>
{
if (s is UdpSession ss) _broadcasts.TryRemove(ss.Remote.Port, out _);
};
}
}
else
{
// 会话集合拒绝(端点重复,或会话已被业务在 NewSession 中释放)。新建的会话没能入集合,
// 就没有任何持有者,必须当场释放:它不会随集合的超时清理被回收,也不该当作“集合里的会话”返回给调用方。
// 释放后回退到集合里已有的会话,调用方拿到的仍是该端点的有效会话(与 TcpServer.OnAccept 的失败分支同型)
session.TryDispose();
// 集合里取不到该端点的会话,说明新建的会话是被业务在 NewSession 中主动释放的(端点重复时这里必能取到)。
// 此处返回这个已释放实例是有意为之:调用方(ProcessReceive 的 OnPreReceive)只用它做判空,
// 真正交付数据的是随后的 OnReceive,它会重新创建该端点的会话并加入集合。
// 若改成抛异常,本轮数据报会被直接丢弃、也不再为该端点建会话(既有用例锁定此行为)
session = sessions.Get(remoteEP) ?? session;
}
}
return session;
}
private void CloseAllSession()
{
var sessions = _Sessions;
if (sessions != null)
{
if (sessions.Count > 0)
{
WriteLog("准备释放会话{0}个!", sessions.Count);
sessions.CloseAll(nameof(CloseAllSession));
sessions.Dispose();
sessions.Clear();
}
}
}
#endregion
#region IServer接口
void IServer.Start() => Open();
void IServer.Stop(String? reason) => Close(reason ?? "Stop");
#endregion
#region 辅助
/// <summary>已重载。</summary>
/// <returns></returns>
public override String ToString()
{
var ss = Sessions;
if (ss != null && ss.Count > 0)
return $"{Local} [{ss.Count}]";
else
return Local.ToString();
}
#endregion
}
/// <summary>Udp扩展</summary>
public static class UdpHelper
{
/// <summary>发送数据流</summary>
/// <param name="udp"></param>
/// <param name="stream"></param>
/// <param name="remoteEP"></param>
/// <returns>返回自身,用于链式写法</returns>
public static UdpClient Send(this UdpClient udp, Stream stream, IPEndPoint? remoteEP = null)
{
Int64 total = 0;
using var buffer = Pool.Rent(1472);
while (true)
{
var n = stream.Read(buffer, 0, buffer.Length);
if (n <= 0) break;
udp.Send(buffer, n, remoteEP);
total += n;
if (n < buffer.Length) break;
}
return udp;
}
/// <summary>向指定目的地发送信息</summary>
/// <param name="udp"></param>
/// <param name="buffer">缓冲区</param>
/// <param name="remoteEP"></param>
/// <returns>返回自身,用于链式写法</returns>
public static UdpClient Send(this UdpClient udp, Byte[] buffer, IPEndPoint? remoteEP = null)
{
udp.Send(buffer, buffer.Length, remoteEP);
return udp;
}
/// <summary>向指定目的地发送信息</summary>
/// <param name="udp"></param>
/// <param name="message"></param>
/// <param name="encoding">文本编码,默认null表示UTF-8编码</param>
/// <param name="remoteEP"></param>
/// <returns>返回自身,用于链式写法</returns>
public static UdpClient Send(this UdpClient udp, String message, Encoding? encoding = null, IPEndPoint? remoteEP = null)
{
if (encoding == null)
Send(udp, Encoding.UTF8.GetBytes(message), remoteEP);
else
Send(udp, encoding.GetBytes(message), remoteEP);
return udp;
}
/// <summary>广播数据包</summary>
/// <param name="udp"></param>
/// <param name="buffer">缓冲区</param>
/// <param name="port"></param>
public static UdpClient Broadcast(this UdpClient udp, Byte[] buffer, Int32 port)
{
if (udp.Client != null && udp.Client.LocalEndPoint != null)
{
if (udp.Client.LocalEndPoint is IPEndPoint ip && !ip.Address.IsIPv4())
throw new NotSupportedException("IPv6 does not support broadcasting!");
}
if (!udp.EnableBroadcast) udp.EnableBroadcast = true;
udp.Send(buffer, buffer.Length, new IPEndPoint(IPAddress.Broadcast, port));
return udp;
}
/// <summary>广播字符串</summary>
/// <param name="udp"></param>
/// <param name="message"></param>
/// <param name="port"></param>
public static UdpClient Broadcast(this UdpClient udp, String message, Int32 port)
{
var buffer = Encoding.UTF8.GetBytes(message);
return Broadcast(udp, buffer, port);
}
/// <summary>接收字符串</summary>
/// <param name="udp"></param>
/// <param name="encoding">文本编码,默认null表示UTF-8编码</param>
/// <returns></returns>
public static String ReceiveString(this UdpClient udp, Encoding? encoding = null)
{
IPEndPoint? ep = null;
var buffer = udp.Receive(ref ep);
if (buffer == null || buffer.Length <= 0) return String.Empty;
encoding ??= Encoding.UTF8;
return encoding.GetString(buffer);
}
}
|