using System.Buffers;
using System.ComponentModel;
using System.Net;
using System.Runtime.ExceptionServices;
using NewLife;
using NewLife.Data;
using NewLife.Messaging;
using NewLife.Net;
using Xunit;
namespace XUnitTest.Net;
/// <summary>UdpSession ç©ºæ•°æ®æŠ¥çº¦å®šçš„åˆ¤å®šæ—¶æœºæµ‹è¯•</summary>
/// <remarks>
/// çº¦å®šï¼šæ”¶åˆ°ç©ºæ•°æ®æŠ¥æ—¶ç»“æŸä¼šè¯ã€‚判定必须以事件链之å‰çš„åŽŸå§‹æ•°æ®æŠ¥ä¸ºå‡†â€”â€”
/// 事件内业务按契约消费/é‡Šæ”¾æœ¬è½®å¥æŸ„åŽ <c>Packet.Length</c> 归零,
/// 事件之åŽå†åˆ¤ä¼šæŠŠæ£å¸¸æ•°æ®æŠ¥è¯¯åˆ¤ä¸ºç©ºåŒ…而销æ¯ä¼šè¯ã€‚
/// </remarks>
public class UdpSessionTests
{
#region 工具
private static readonly IPEndPoint _remote = new(IPAddress.Loopback, 12345);
private static UdpSession NewSession(UdpServer server) => new(server, null, _remote);
private static ReceivedEventArgs NewArgs(IPacket? pk) => new()
{
Local = IPAddress.Loopback,
Remote = _remote,
Packet = pk,
};
/// <summary>æž„é€ å¸¦æ‹¥æœ‰æƒçš„æ•°æ®æŠ¥å¥æŸ„:缓冲å–è‡ªæ± ï¼Œä¿è¯é‡Šæ”¾åŽåŽŸæ ·å½’è¿˜</summary>
private static OwnerPacket NewPacket(Int32 length) => new(ArrayPool<Byte>.Shared.Rent(16), 0, length, true);
/// <summary>测试用åè®®ï¼šæ•´ä¸ªæ•°æ®æŠ¥å³ä¸€æ¡æ¶ˆæ¯ï¼ˆæ— 头部å—节)</summary>
private sealed class DatagramCodec : IMessageCodec
{
public ParseResult? TryParse(ReadOnlySequence<Byte> buffer)
{
if (buffer.Length == 0) return null;
return new ParseResult { Message = new Message(), HeaderSize = 0, BodyLength = buffer.Length };
}
public IOwnerPacket? Build(IMessage message) => throw new NotSupportedException("测试åè®®ä»…éªŒè¯æŽ¥æ”¶è·¯å¾„");
public IOwnerPacket BuildHeader(IMessage message, Int64 bodyLength) => throw new NotSupportedException("测试åè®®ä»…éªŒè¯æŽ¥æ”¶è·¯å¾„");
}
#endregion
[Theory(DisplayName = "UDPäº‹ä»¶å†…æ¶ˆè´¹æˆ–ç½®ç©ºå¥æŸ„_ä¸å¾—è¯¯åˆ¤ä¸ºç©ºæ•°æ®æŠ¥åœä¼šè¯")]
[InlineData("dispose")]
[InlineData("null")]
public void ConsumedInHandler_NotTreatedAsEmpty(String mode)
{
using var server = new UdpServer();
using var session = NewSession(server);
var handled = 0;
session.Received += (s, e) =>
{
handled++;
// 契约å…è®¸çš„ä¸¤ç±»æ”¹å†™ï¼šæ¶ˆè´¹æœ¬è½®å¥æŸ„(交å‘é€é“¾è·¯æˆ–显å¼é‡Šæ”¾ï¼‰ã€ç½®ç©ºåšæ ‡è®°
if (mode == "dispose")
e.Packet?.TryDispose();
else
e.Packet = null;
};
var pk = NewPacket(3);
try
{
session.OnReceive(NewArgs(pk));
}
finally
{
pk.TryDispose();
}
Assert.Equal(1, handled);
// éžç©ºæ•°æ®æŠ¥ä¸å¾—è§¦å‘ Stop+Dispose:会è¯å˜æ´»ã€æœåŠ¡å™¨å¼•ç”¨ä¿ç•™
Assert.False(session.Disposed);
Assert.Same(server, session.Server);
}
[Fact(DisplayName = "UDPç©ºæ•°æ®æŠ¥_结æŸä¼šè¯")]
public void EmptyDatagram_StopsSession()
{
using var server = new UdpServer();
using var session = NewSession(server);
var handled = 0;
session.Received += (s, e) => handled++;
var pk = NewPacket(0);
try
{
session.OnReceive(NewArgs(pk));
}
finally
{
pk.TryDispose();
}
// äº‹ä»¶ç…§æ—§æŠ›å‡ºï¼ˆä¸šåŠ¡å¯æ„ŸçŸ¥ï¼‰ï¼Œä½†ä¼šè¯æŒ‰çº¦å®šç»“æŸ
Assert.Equal(1, handled);
Assert.True(session.Disposed);
Assert.Null(session.Server);
}
/// <summary>并行派å‘ä¸‹çš„å¹¶å‘æ”¶åŒ…:消æ¯ä¸å¾—滞留ã€å¹¶å‘æ§½ä½ä¸å¾—溢出</summary>
/// <remarks>
/// 首访竞æ€çª—壿žçª„ï¼ˆè¯»å—æ®µåˆ°èµ‹å€¼åªæœ‰å‡ åçº³ç§’ï¼‰ï¼Œæ— æ³•ç¡®å®šå¤çŽ°â€”â€”å®žæµ‹æ—§å®žçŽ°è·‘æœ¬ç”¨ä¾‹åŒæ ·é€šè¿‡ã€‚
/// 本用例守护的是并行派å‘å¥‘çº¦ï¼šå¹¶å‘æ³¨å…¥çš„æ•°æ®æŠ¥å¿…须全部处ç†å®Œï¼Œä¸”ä¸å¾—å‡ºçŽ°æ§½ä½æº¢å‡º
/// (溢出会从 fire-and-forget ä»»åŠ¡é€ƒé€¸ï¼Œæ— äººè§‚å¯Ÿã€æ— 日志,åªèƒ½é 首轮异常通知检出)。
/// </remarks>
[Fact(DisplayName = "UDP并行派å‘_å¹¶å‘é¦–è®¿ä¸Žå¹¶å‘æ”¶åŒ…_消æ¯ä¸ä¸¢ä¸”æ— æ§½ä½æº¢å‡º")]
public void ParallelDispatch_ConcurrentFirstAccess_NoSlotOverflow()
{
using var server = new UdpServer();
using var session = NewSession(server);
session.Protocol = new DatagramCodec();
session.MaxConcurrency = 4;
const Int32 workers = 8;
const Int32 rounds = 8;
const Int32 count = workers * rounds;
using var done = new CountdownEvent(count);
session.Received += (s, e) => done.Signal();
// å¹¶å‘ä¿¡å·é‡çš„首访竞æ€ä¼šè®©â€œç‰å¾…的实例â€ä¸Žâ€œå½’还的实例â€ä¸åŒï¼Œæ§½ä½è´¦ç›®å¤±è¡¡ã€‚
// 溢出异常从 fire-and-forget ä»»åŠ¡é€ƒé€¸ï¼ˆæ— äººè§‚å¯Ÿã€æ— 日志),åªèƒ½é 首轮异常通知æ•获
var overflow = 0;
EventHandler<FirstChanceExceptionEventArgs> handler = (s, e) =>
{
if (e.Exception is SemaphoreFullException) Interlocked.Increment(ref overflow);
};
AppDomain.CurrentDomain.FirstChanceException += handler;
try
{
// Barrier è®©æ‰€æœ‰çº¿ç¨‹åŒæ—¶é¦–访并å‘ä¿¡å·é‡ï¼Œå†å¹¶å‘æ³¨å…¥æ•°æ®æŠ¥
using var barrier = new Barrier(workers);
var threads = new Thread[workers];
for (var i = 0; i < workers; i++)
{
threads[i] = new Thread(() =>
{
barrier.SignalAndWait();
for (var j = 0; j < rounds; j++)
{
var pk = NewPacket(3);
try
{
session.OnReceive(NewArgs(pk));
}
finally
{
pk.TryDispose();
}
}
})
{ IsBackground = true };
}
foreach (var th in threads) th.Start();
foreach (var th in threads) Assert.True(th.Join(10_000), "并呿”¶åŒ…未在超时内完æˆ");
Assert.True(done.Wait(TimeSpan.FromSeconds(10)), $"å¹¶è¡Œå¤„ç†æœªåœ¨è¶…时内完æˆï¼Œå·²å¤„ç† {count - done.CurrentCount}/{count}");
}
finally
{
AppDomain.CurrentDomain.FirstChanceException -= handler;
}
Assert.Equal(0, overflow);
}
}
|