using System.Diagnostics;
using System.Globalization;
using System.Net;
using System.Net.Sockets;
using NewLife;
using NewLife.Collections;
using NewLife.Net;
namespace NetLoadTest;
/// <summary>裸 Socket 层压测程åºï¼šecho æœåŠ¡ç«¯ï¼ˆå¯ç‹¬ç«‹è¿›ç¨‹ï¼‰+ N 客户端,测é‡åžåã€å¾€è¿”延迟与内å˜åˆ†é…</summary>
/// <remarks>
/// <para>用法:dotnet run --project Benchmark/NetLoadTest -c Release -- [--mode pipeline|roundtrip] [--clients 4] [--size 1024] [--seconds 10] [--warmup 2] [--recvmode sync|asyncpull|event] [--udp] [--frame 24]</para>
/// <para>pipeline:客户端æŒç»å‘é€ä¸ç‰å¾…回显,测åžå上é™ï¼ˆMB/sã€msg/s);roundtrip:é€åŒ…往返,测 P50/P95/P99 延迟;--oneway:åªå‘ä¸è¯»ï¼Œé…åˆ --server 分离进程测æœåŠ¡ç«¯çº¯æŽ¥æ”¶åžå。</para>
/// <para>--recvmode:客户端接收模å¼ã€‚sync=åŒæ¥æ‹‰å– Receive(阻塞ç‰å¾…);asyncpull=å¼‚æ¥æ‹‰å– ReceiveAsyncï¼›event=事件接收(AutoReceive=true + Received 回调,与æœåŠ¡å™¨æŽ¥æ”¶æ¨¡åž‹åŒæž„,零阻塞线程)。</para>
/// <para>--frame:应用层帧大å°ï¼Œå‘é€ç¼“冲对é½åˆ°å¸§æ•´å€æ•°æ¨¡æ‹Ÿç²˜åŒ…,åžå折算为逻辑帧å£å¾„ï¼ˆå¯¹æ ‡åŽ†å² 1.4 亿 pkt/s)。</para>
/// <para>--server:独立æœåŠ¡ç«¯è¿›ç¨‹ï¼Œæ”¶åˆ°æ•°æ®åŽé™é»˜ 3 秒输出 STEADY 稳æ€ä¸ä½æ•°å¹¶é€€å‡ºï¼ˆå«æœåŠ¡ç«¯åˆ†é… B/msg);--remote 客户端连接独立æœåŠ¡ç«¯ã€‚</para>
/// <para>--udpconnect:UDP 客户端将 Socket Connect åˆ°ç›®æ ‡ç«¯ï¼ˆå®žéªŒï¼‰ï¼Œå‘é€èµ°å·²è¿žæŽ¥è·¯å¾„(零分é…)。</para>
/// <para>åˆ†é…æ•°æ®æ¥è‡ªè¿›ç¨‹çº§ GC.GetTotalAllocatedBytes(分离进程时仅本进程开销)。末尾输出 SUMMARY:{json} 机器å¯è¯»æ±‡æ€»ï¼Œä¾›è·‘批脚本解æžã€‚</para>
/// </remarks>
static class Program
{
private static Int64 _serverBytes;
private static Int64 _sentBytes;
public static void Main(String[] args)
{
// 统一 UTF-8 输出,ä¿è¯é‡å®šå‘æ—¥å¿—ï¼ˆè·‘æ‰¹è„šæœ¬é‡‡é›†ï¼‰ä¸æ–‡ä¸ä¹±ç
Console.OutputEncoding = System.Text.Encoding.UTF8;
// 高并å‘下大é‡é˜»å¡žå¼æ”¶å‘ä»»åŠ¡ä¼šè§¦çº¿ç¨‹æ± æ³¨å…¥é™é€Ÿï¼ˆ~1线程/ç§’ï¼‰ï¼Œé¢„çƒæœ€å°çº¿ç¨‹é¿å…秒级毛刺干扰测é‡
ThreadPool.SetMinThreads(512, 512);
var mode = GetArg(args, "--mode") ?? "pipeline";
var clients = GetInt(args, "--clients", 4);
var size = GetInt(args, "--size", 1024);
var seconds = GetInt(args, "--seconds", 10);
var warmup = GetInt(args, "--warmup", 2);
var udp = args.Contains("--udp");
var serverOnly = args.Contains("--server");
var remote = GetArg(args, "--remote");
var port = GetInt(args, "--port", 7800);
var oneway = args.Contains("--oneway");
var frame = GetInt(args, "--frame", 0);
var udpConnect = args.Contains("--udpconnect");
var recvMode = (GetArg(args, "--recvmode") ?? "sync").ToLowerInvariant();
if (recvMode is not ("sync" or "asyncpull" or "event"))
{
Console.WriteLine($"æœªçŸ¥æŽ¥æ”¶æ¨¡å¼ --recvmode {recvMode},å¯é€‰ sync|asyncpull|event");
return;
}
// UDP å•包上é™ï¼šIP 首部 20 + UDP 首部 8 åŽæœ€å¤§è´Ÿè½½ 65507ï¼Œè¶…å‡ºæ—¶å†…æ ¸ç›´æŽ¥æ‹’ç»ï¼ˆSend 返回 -1)
if (udp && size > 65507) size = 65507;
// --frame:应用层帧大å°ã€‚å‘é€ç¼“冲对é½åˆ°å¸§æ•´å€æ•°ï¼Œæ¨¡æ‹Ÿâ€œå¤§é‡å°å¸§ç²˜æˆå¤§åŒ…â€çš„å议场景,
// 统计å£å¾„折算为逻辑帧åžåï¼ˆå¯¹æ ‡åŽ†å² 23.4Gbps ÷ 24B = 1.4 亿 pkt/s 记录)
if (frame > 0) size = Math.Max(frame, size / frame * frame);
// å†…æ ¸æŽ¥æ”¶ç¼“å†²ï¼ˆTCP 大报文åžå上é™ï¼‰ï¼š0 è¡¨ç¤ºä¿æŒç³»ç»Ÿé»˜è®¤
var rcvbuf = GetInt(args, "--rcvbuf", 0);
if (rcvbuf > 0) SocketSetting.Current.ReceiveBufferSize = rcvbuf;
var roundtrip = mode.Equals("roundtrip", StringComparison.OrdinalIgnoreCase);
Console.WriteLine("=== 裸Socket压测(echo æœåŠ¡ç«¯ï¼‰===");
Console.WriteLine($"æ¨¡å¼ : {(roundtrip ? "é€åŒ…往返(延迟)" : oneway ? "å•å‘上行(åžå)" : "æµæ°´çº¿ï¼ˆåžå)")}");
Console.WriteLine($"åè®® : {(udp ? "UDP" : "TCP")}");
Console.WriteLine($"åŒ…å¤§å° : {size:N0} B{(frame > 0 ? $"ï¼ˆå« {size / frame:N0} 个 {frame} B 逻辑帧,粘包å£å¾„)" : "")}");
Console.WriteLine($"客户端 : {clients} 接收模å¼: {recvMode}{(udpConnect ? " UDP-Connect实验" : "")}");
Console.WriteLine($"预çƒ/窗å£: {warmup} s / {seconds} s");
Console.WriteLine();
// ===== æœåŠ¡ç«¯ï¼ˆecho + 计数) =====
// å¢žå¤§ä¼šè¯æŽ¥æ”¶ç¼“å†²ï¼Œé™ä½Žå¤§åŒ…分段的系统调用开销
SocketSetting.Current.BufferSize = Math.Max(64 * 1024, size);
NetServer? server = null;
if (remote == null || serverOnly)
{
server = new NetServer
{
Port = serverOnly ? port : 0,
ProtocolType = udp ? NetType.Udp : NetType.Tcp,
AddressFamily = AddressFamily.InterNetwork,
};
server.Received += (s, e) =>
{
var pk = e.Packet;
if (pk == null || pk.Length <= 0) return;
Interlocked.Add(ref _serverBytes, pk.Length);
// å•å‘上行模å¼åªè®¡æ•°ä¸å›žå‘ï¼Œç”¨äºŽæµ‹é‡æœåŠ¡ç«¯çº¯æŽ¥æ”¶åžå(对é½åކå²å£å¾„)
if (oneway) return;
// TCP æœåŠ¡ç«¯ sender 为 NetSession(INetSession);UDP æœåŠ¡ç«¯ sender 为 UdpSession(ISocketSession),两者接å£ä¸åŒ
if (s is INetSession ns) ns.Send(pk);
else if (s is UdpSession us) us.Send(pk);
};
server.Start();
}
if (serverOnly)
{
// 独立æœåŠ¡ç«¯è¿›ç¨‹ï¼šæ‰“å°å°±ç»ªä¸Žæ¯ç§’接收速率;收到数æ®åŽè¿žç»é™é»˜ 3 秒视为客户端结æŸï¼Œè¾“出稳æ€ä¸ä½æ•°å¹¶é€€å‡º
Console.WriteLine($"READY {(udp ? "udp" : "tcp")}://0.0.0.0:{server!.Port}");
var last = 0L;
var rates = new List<Double>();
var idle = 0;
// æœåŠ¡ç«¯åˆ†é…统计:接收路径(å«å›žæ˜¾å‘é€ï¼‰çš„æ‰˜ç®¡åˆ†é…总é‡ï¼ŒæŒ‰æŽ¥æ”¶æ¶ˆæ¯æ•°æŠ˜ç®— B/msg
var srvAlloc0 = GC.GetTotalAllocatedBytes(false);
var srvPause0 = GC.GetTotalPauseDuration();
var srvGc0 = GC.CollectionCount(0);
var srvGc1 = GC.CollectionCount(1);
var srvGc2 = GC.CollectionCount(2);
while (true)
{
Thread.Sleep(1000);
var now = Interlocked.Read(ref _serverBytes);
var delta = now - last;
var rate = delta / 1024.0 / 1024.0;
if (delta > 0)
{
rates.Add(rate);
idle = 0;
}
else
{
idle++;
}
if (frame > 0)
{
// 粘包å£å¾„:åžå按逻辑帧折算(接收å—节 ÷ 帧大å°ï¼‰
Console.WriteLine($"[server] {rate:N1} MB/s 累计 {now / 1024.0 / 1024.0:N1} MB 帧率 {delta / (Double)frame / 1_000_000:N2} M帧/s");
}
else
{
Console.WriteLine($"[server] {rate:N1} MB/s 累计 {now / 1024.0 / 1024.0:N1} MB");
}
last = now;
if (idle >= 3 && rates.Count > 0)
{
// ä¸¢æŽ‰å‰ 2 ç§’çˆ¬å‡æ®µï¼ˆè¿žæŽ¥å»ºç«‹/æ…¢å¯åŠ¨ï¼‰ï¼Œå–å‰©ä½™ç§’é‡‡æ ·ä¸ä½æ•°ä½œä¸ºç¨³æ€
var window = rates.Count > 6 ? rates.GetRange(2, rates.Count - 2) : rates;
var steady = Median(window);
var msgs = now / size;
var srvAllocPerMsg = msgs > 0 ? (GC.GetTotalAllocatedBytes(false) - srvAlloc0) / (Double)msgs : 0;
var line = $"[server] STEADY {steady:N1} MB/s Gbps={steady * 0.008388608:N2} pkt={steady * 1048576 / size:N0} pkt/s";
if (frame > 0) line += $" frame={steady * 1048576 / frame:N0} frame/s";
line += $" alloc={srvAllocPerMsg:N2} B/msg gc={GC.CollectionCount(0) - srvGc0}/{GC.CollectionCount(1) - srvGc1}/{GC.CollectionCount(2) - srvGc2} æš‚åœ={(GC.GetTotalPauseDuration() - srvPause0).TotalMilliseconds:N1} ms ç§’é‡‡æ ·={window.Count}";
Console.WriteLine(line);
break;
}
// 从未收到任何数æ®ï¼ˆå®¢æˆ·ç«¯å¼‚常或包超é™è¢«å†…æ ¸æ‹’ç»ï¼‰ï¼š10 ç§’åŽé€€å‡ºï¼Œé¿å…跑批脚本挂ç‰
if (idle >= 10 && rates.Count == 0)
{
Console.WriteLine("[server] STEADY 0.0 MB/s Gbps=0 pkt=0 pkt/s alloc=0 gc=0/0/0 ç§’é‡‡æ ·=0");
break;
}
}
server.Dispose();
return;
}
// ===== 客户端 =====
var payload = new Byte[size];
Random.Shared.NextBytes(payload);
// event 模å¼èµ°æŽ¥æ”¶çŽ¯äº‹ä»¶æŽ¨é€ï¼ˆAutoReceive=true);sync/asyncpull ä¸ºæ‹‰å–æ¨¡å¼ï¼ˆæ‰“å¼€å‰å¿…须关自动接收)
var autoReceive = recvMode == "event";
var conns = new List<ISocketClient>();
for (var i = 0; i < clients; i++)
{
var hostPort = remote ?? $"127.0.0.1:{server!.Port}";
ISocketClient conn;
if (udp)
{
conn = new UdpServer
{
Remote = new NetUri($"udp://{hostPort}"),
AutoReceive = autoReceive,
// UDP æ— æµæŽ§å¯èƒ½ä¸¢åŒ…ï¼šçŸæŽ¥æ”¶è¶…æ—¶ä¾¿äºŽæŽ’æ°´é˜¶æ®µè¯†åˆ«â€œä¸å†æœ‰æ•°æ®â€
Timeout = 300,
};
}
else
{
conn = new TcpSession
{
Remote = new NetUri($"tcp://{hostPort}"),
AutoReceive = autoReceive,
BufferSize = Math.Max(64 * 1024, size),
Timeout = 30_000,
};
}
conn.Open();
// UDP Connect 化实验:连接åŽå‘é€èµ°å·²è¿žæŽ¥è·¯å¾„(零分é…ã€çœ SendTo åºåˆ—化开销)
if (udp && udpConnect && conn is SessionBase sb && sb.Client is Socket usk)
{
var idx = hostPort.LastIndexOf(':');
try
{
usk.Connect(new IPEndPoint(IPAddress.Parse(hostPort[..idx]), Int32.Parse(hostPort[(idx + 1)..])));
}
catch (Exception ex)
{
Console.WriteLine($"UDP Connect 失败:{ex.Message}");
}
}
conns.Add(conn);
}
// 预çƒï¼šè®©è¿žæŽ¥ä¸ŽæŽ¥æ”¶çŽ¯è¿›å…¥ç¨³å®šçŠ¶æ€
Thread.Sleep(warmup * 1000);
// é¢„çƒæµé‡ï¼šçŸä¿ƒæ”¶å‘è§¦å‘æœåŠ¡ç«¯å›žæ˜¾/接收链路与客户端接收链路的首次 JIT ä¸Žæ± åŒ–åˆå§‹åŒ–。
// æ¯ä¸ªåœºæ™¯éƒ½æ˜¯æ–°èµ·çš„æœåŠ¡ç«¯è¿›ç¨‹ï¼Œè‹¥ä¸åšè¿™è½®é¢„çƒï¼Œé¦–ä¸ªæµ‹é‡æ ·æœ¬ä¼šæŠŠæœåŠ¡ç«¯å†·å¯åŠ¨æˆæœ¬ï¼ˆç™¾æ¯«ç§’级)计入尾延迟。
PrimeTrafficAsync(oneway, udp, remote, server, size, recvMode, payload).GetAwaiter().GetResult();
var alloc0 = GC.GetTotalAllocatedBytes(false);
var pause0 = GC.GetTotalPauseDuration();
var gen0 = GC.CollectionCount(0);
var gen1 = GC.CollectionCount(1);
var gen2 = GC.CollectionCount(2);
var bytes0 = Interlocked.Read(ref _serverBytes);
var sent0 = Interlocked.Read(ref _sentBytes);
var sw = Stopwatch.StartNew();
List<Double>? samples = null;
Int64 sent = 0, received = 0;
if (roundtrip)
samples = RunRoundTrip(conns, payload, size, seconds, udp, recvMode);
else if (oneway)
RunOneWay(conns, payload, size, seconds);
else
(sent, received) = RunPipeline(conns, payload, size, seconds, udp, recvMode);
sw.Stop();
var sent1 = Interlocked.Read(ref _sentBytes);
var alloc1 = GC.GetTotalAllocatedBytes(false);
// ===== 统计 =====
// å•å‘上行按固定å‘é€çª—计时(排水阶段ä¸è®¡å…¥çª—å£ï¼‰ï¼›å…¶ä½™æŒ‰å®žé™…墙钟
var elapsed = oneway ? (Double)seconds : sw.Elapsed.TotalSeconds;
var sentBytes = sent1 - sent0;
// pipeline åžåå–客户端收到的回显å—节;oneway æ— å›žè¯»å–å‘é€å—节
var totalBytes = oneway ? sentBytes : received;
var msgCount = roundtrip ? samples?.Count ?? 0 : totalBytes / size;
var allocPerMsg = msgCount > 0 ? (alloc1 - alloc0) / (Double)msgCount : 0;
Console.WriteLine();
Console.WriteLine("------- 结果 -------");
if (!roundtrip)
{
var mbps = totalBytes / elapsed / (1024.0 * 1024.0);
Console.WriteLine($"åžå : {totalBytes / size / elapsed:N0} msg/s | {mbps:N1} MB/s");
if (frame > 0)
{
// 粘包å£å¾„:按逻辑帧折算åžåï¼Œå¯¹æ ‡åŽ†å²â€œ1.4 亿 pkt/sâ€ï¼ˆ23.4Gbps ÷ 24B)
var totalFrames = totalBytes / frame;
Console.WriteLine($"帧åžå : {totalFrames / elapsed:N0} frame/sï¼ˆå¸§å¤§å° {frame} B,æ¯å¤§åŒ… {size / frame:N0} 帧)");
}
Console.WriteLine($"{(oneway ? "å‘é€" : "回显")} : {totalBytes / size:N0} 包 / {totalBytes:N0} Bï¼ˆçª—å£ {elapsed:F2} s)");
}
Console.WriteLine($"åˆ†é… : {allocPerMsg:N1} B/msg | çª—å£æ€»åˆ†é… {(alloc1 - alloc0) / (1024.0 * 1024.0):N1} MB");
Console.WriteLine($"GC : Gen0 +{GC.CollectionCount(0) - gen0} Gen1 +{GC.CollectionCount(1) - gen1} Gen2 +{GC.CollectionCount(2) - gen2}");
Console.WriteLine($"GCæš‚åœ : {(GC.GetTotalPauseDuration() - pause0).TotalMilliseconds:N2} msï¼ˆçª—å£ {elapsed:F2} s)");
// 机器å¯è¯»æ±‡æ€»ï¼ˆè·‘批脚本按 SUMMARY: å‰ç¼€è§£æž JSON)
var summary = new Dictionary<String, Object?>
{
["mode"] = roundtrip ? "roundtrip" : oneway ? "oneway" : "pipeline",
["protocol"] = udp ? "udp" : "tcp",
["recvMode"] = recvMode,
["clients"] = clients,
["size"] = size,
["frame"] = frame,
["elapsed"] = Math.Round(elapsed, 3),
["sentBytes"] = sentBytes,
["allocPerMsg"] = Math.Round(allocPerMsg, 2),
["gc0"] = GC.CollectionCount(0) - gen0,
["gc1"] = GC.CollectionCount(1) - gen1,
["gc2"] = GC.CollectionCount(2) - gen2,
["gcPauseMs"] = Math.Round((GC.GetTotalPauseDuration() - pause0).TotalMilliseconds, 2),
};
if (roundtrip)
{
summary["latencyCount"] = samples?.Count ?? 0;
if (samples is { Count: > 0 })
{
summary["p50"] = Math.Round(Percentile(samples, 0.50), 3);
summary["p95"] = Math.Round(Percentile(samples, 0.95), 3);
summary["p99"] = Math.Round(Percentile(samples, 0.99), 3);
summary["avg"] = Math.Round(samples.Average(), 3);
summary["max"] = Math.Round(samples[^1], 3);
}
}
else
{
summary["recvBytes"] = totalBytes;
summary["msgsPerSec"] = Math.Round(totalBytes / size / elapsed, 1);
summary["mbps"] = Math.Round(totalBytes / elapsed / (1024.0 * 1024.0), 2);
if (frame > 0) summary["framesPerSec"] = Math.Round(totalBytes / frame / elapsed, 1);
if (!oneway)
{
var sentMsgs = sentBytes / size;
var recvMsgs = totalBytes / size;
summary["integrityDiff"] = sentMsgs - recvMsgs;
if (udp && sentMsgs > 0) summary["lossPct"] = Math.Round((sentMsgs - recvMsgs) * 100.0 / sentMsgs, 2);
}
}
Console.WriteLine("SUMMARY:" + ToJson(summary));
foreach (var conn in conns) conn.Dispose();
server?.Dispose();
}
/// <summary>å•å‘上行:仅å‘é€ä¸å›žè¯»ï¼Œæµ‹æœåŠ¡ç«¯çº¯æŽ¥æ”¶åžå(é…åˆ --server 分离进程消除 CPU 共享)</summary>
private static void RunOneWay(List<ISocketClient> conns, Byte[] payload, Int32 size, Int32 seconds)
{
var tasks = new List<Task>();
for (var i = 0; i < conns.Count; i++)
{
var conn = conns[i];
tasks.Add(Task.Run(() =>
{
var n = 0L;
var batch = 0L;
// å„客户端独立时间窗,到点å³åœï¼ˆé˜»å¡žä¸çš„ Send 返回åŽç«‹å³é€€å‡ºï¼‰
var deadline = Stopwatch.GetTimestamp() + (Int64)(seconds * Stopwatch.Frequency);
var errors = 0;
while (Stopwatch.GetTimestamp() < deadline)
{
// å‘é€å¤±è´¥ï¼ˆå¦‚ UDP 包超é™è¿”回 -1)ä¸è®¡å…¥å‘é€é‡ï¼Œè¿žç»å¤±è´¥åˆ™æå‰é€€å‡ºé˜²æ»è½¬
if (conn.Send(payload) <= 0)
{
if (++errors > 1000) break;
continue;
}
errors = 0;
n++;
batch += payload.Length;
if ((n & 0xFF) == 0)
{
Interlocked.Add(ref _sentBytes, batch);
batch = 0;
}
}
Interlocked.Add(ref _sentBytes, batch);
}));
}
// å‘é€çª—结æŸåŽä»æœ‰å°‘é‡åœ¨é€”:TCP æµæŽ§ä¸‹ç¼“å†²æ»¡æ—¶ä¼šç¨æ™šè¿”回,宽é™ç‰å¾…。
// æœåŠ¡ç«¯é¥±å’Œåœºæ™¯ä¸ Send å¯èƒ½é•¿æœŸé˜»å¡žäºŽé›¶çª—å£ï¼Œå®½é™ 10 ç§’åŽæ”¾å¼ƒç‰å¾…:
// å‘é€ç»Ÿè®¡ç›´æŽ¥è¯»å…¨å±€ç´¯è®¡ï¼ˆ_sentBytes),ä¸ä¾èµ–阻塞任务是å¦è¿”回
Task.WaitAll(tasks.ToArray(), TimeSpan.FromSeconds(10));
var totalBytes = Interlocked.Read(ref _sentBytes);
Console.WriteLine($"完整å‘é€ : {totalBytes / size:N0} 包 / {totalBytes:N0} Bï¼ˆæ— å›žè¯»ï¼ŒæŽ¥æ”¶çœŸå€¼ä»¥æœåŠ¡ç«¯è®¡æ•°ä¸ºå‡†ï¼‰");
}
/// <summary>é¢„çƒæµé‡ï¼šè§¦å‘æœåŠ¡ç«¯æ”¶å‘链路与客户端接收链路的首次 JIT ä¸Žæ± åŒ–åˆå§‹åŒ–,é¿å…冷å¯åŠ¨æˆæœ¬è®¡å…¥æµ‹é‡</summary>
private static async Task PrimeTrafficAsync(Boolean oneway, Boolean udp, String? remote, NetServer? server, Int32 size, String recvMode, Byte[] payload)
{
var hostPort = remote ?? $"127.0.0.1:{server!.Port}";
// æ€»é¢„çƒæµé‡æŽ§åˆ¶åœ¨ 4MB é‡çº§ï¼šæ¬¡æ•°å¤Ÿè§¦å‘ JIT,大包åˆä¸è‡³äºŽæŠŠå†…æ ¸ç¼“å†²å µæ»¡
var count = Math.Clamp(4 * 1024 * 1024 / size, 4, 64);
// å•呿¨¡å¼æœåŠ¡ç«¯ä¸å›žæ˜¾ï¼šåªå‘䏿”¶ï¼Œè§¦å‘æœåŠ¡ç«¯æŽ¥æ”¶é“¾è·¯ JIT
if (oneway)
{
using var pc = CreateClient(udp, hostPort, false, size);
pc.Open();
for (var i = 0; i < count; i++) pc.Send(payload);
Thread.Sleep(200);
return;
}
// 事件模å¼ï¼šå›žæ‰§èµ°æŽ¥æ”¶çŽ¯å›žè°ƒï¼Œé¢„çƒå®¢æˆ·ç«¯äº‹ä»¶é“¾è·¯ï¼ˆä¸Žæµ‹é‡æ¨¡å¼ä¸€è‡´ï¼‰
if (recvMode == "event")
{
var pc = CreateClient(udp, hostPort, true, size);
try
{
var got = 0L;
using var done = new ManualResetEventSlim();
pc.Received += (s, e) =>
{
var n = e.Packet?.Length ?? 0;
if (n > 0 && Interlocked.Add(ref got, n) >= (Int64)count * size) done.Set();
};
pc.Open();
for (var i = 0; i < count; i++) pc.Send(payload);
done.Wait(3000); // UDP 容å¿ä¸¢åŒ…:超时å³è§†ä¸ºé¢„çƒè¶³å¤Ÿ
}
finally
{
pc.Dispose();
}
return;
}
// æ‹‰å–æ¨¡å¼ï¼šä¸²è¡Œå¾€è¿”(TCP å¼‚æ¥æ‹‰å–èµ° ReceiveAsync,与测é‡è·¯å¾„一致;UDP 用çŸè¶…æ—¶åŒæ¥æŽ¥æ”¶ï¼‰
using (var pc = CreateClient(udp, hostPort, false, size))
{
pc.Open();
for (var i = 0; i < count; i++)
{
pc.Send(payload);
var need = size;
while (need > 0)
{
try
{
using var pk = !udp && recvMode == "asyncpull" ? await pc.ReceiveAsync() : pc.Receive();
if (pk == null || pk.Length <= 0) break;
need -= pk.Length;
}
catch (SocketException ex) when (udp && ex.SocketErrorCode == SocketError.TimedOut)
{
// UDP 丢包:跳过该次
break;
}
}
}
}
}
/// <summary>创建指定接收模å¼çš„客户端连接(UDP 客户端用 UdpServer 承载)</summary>
private static ISocketClient CreateClient(Boolean udp, String hostPort, Boolean autoReceive, Int32 size)
=> udp
? new UdpServer
{
Remote = new NetUri($"udp://{hostPort}"),
AutoReceive = autoReceive,
Timeout = 300,
}
: new TcpSession
{
Remote = new NetUri($"tcp://{hostPort}"),
AutoReceive = autoReceive,
BufferSize = Math.Max(64 * 1024, size),
Timeout = 30_000,
};
/// <summary>æµæ°´çº¿æ¨¡å¼ï¼šæŒç»å‘é€ + å¹¶å‘æŽ¥æ”¶å›žæ˜¾ï¼ˆåŒæ¥/å¼‚æ¥æ‹‰å–或事件接收),å‘é€åœæ¢åŽæŽ’æ°´è¯»æ»¡</summary>
private static (Int64 Sent, Int64 Received) RunPipeline(List<ISocketClient> conns, Byte[] payload, Int32 size, Int32 seconds, Boolean tolerateLoss, String recvMode)
{
var runners = conns.Select(c => new PipelineClient(c, size, tolerateLoss, recvMode)).ToArray();
var sendTasks = new List<Task>();
var recvTasks = new List<Task>();
foreach (var r in runners)
{
sendTasks.Add(r.StartSend(payload));
recvTasks.Add(r.StartReceive());
}
Thread.Sleep(seconds * 1000);
foreach (var r in runners) r.Stop();
// å…ˆç‰å‘é€ä»»åŠ¡æ”¶å°¾ï¼šSentPackets 在å‘é€ä»»åŠ¡ finally æ‰å‘布,未收尾就排水会以åå°çš„å‘é€é‡æå‰åˆ¤å®šâ€œå·²æŽ’å¹²â€
if (!Task.WaitAll(sendTasks.ToArray(), TimeSpan.FromSeconds(30)))
Console.WriteLine("è¦å‘Šï¼šéƒ¨åˆ†å‘é€ä»»åŠ¡æœªåŠæ—¶ç»“æŸ");
// æŽ’æ°´ï¼šè½®è¯¢ç‰æŽ¥æ”¶è¯»æ»¡å‘逿€»é‡ï¼ˆUDP ä¸¢åŒ…æˆ–é•¿æ—¶é—´æ— è¿›å±•åˆ™æ”¾å¼ƒï¼‰
var drained = false;
var deadline = Stopwatch.GetTimestamp() + (Int64)(60 * Stopwatch.Frequency);
var lastTotal = -1L;
var stall = 0;
if (tolerateLoss)
{
// UDP æ— æµæŽ§ï¼šçŒåŒ…å¿…ç„¶ä¸¢åŒ…ï¼ŒæŽ’æ°´æ— æ„义;çŸç»“ç®—åŽç›´æŽ¥ç»Ÿè®¡ï¼ˆå‘é€ä¸ŽæŽ¥æ”¶å·®å³ä¸ºä¸¢åŒ…)
Thread.Sleep(500);
}
else
{
while (true)
{
var sent = runners.Sum(r => r.SentPackets);
var recv = runners.Sum(r => r.TotalReceived);
if (sent > 0 && recv >= sent * size)
{
drained = true;
break;
}
if (Stopwatch.GetTimestamp() > deadline)
{
Console.WriteLine("è¦å‘Šï¼šæŽ’水超时,部分回显未读满(å¯èƒ½å˜åœ¨ä¸¢åŒ…)");
break;
}
if (recv == lastTotal)
{
if (++stall >= 50) break;
}
else
{
stall = 0;
lastTotal = recv;
}
Thread.Sleep(100);
}
}
// 未排干时主动关é—,é¿å…接收侧悬挂
if (!drained)
foreach (var c in conns) c.Close("drain-timeout");
// ç‰æŽ¥æ”¶ä»»åŠ¡æ”¶å°¾ï¼ˆæ‹‰å–æ¨¡å¼å¾ªçŽ¯é€€å‡ºï¼›äº‹ä»¶æ¨¡å¼æ— 接收任务)
if (!Task.WaitAll(recvTasks.ToArray(), TimeSpan.FromSeconds(15)))
Console.WriteLine("è¦å‘Šï¼šéƒ¨åˆ†æŽ¥æ”¶ä»»åŠ¡æœªåŠæ—¶ç»“æŸ");
var sentBytes = runners.Sum(r => r.SentPackets) * size;
var recvBytes = runners.Sum(r => r.TotalReceived);
Console.WriteLine($"完整性 : å‘é€ {sentBytes / size:N0} 包 / 接收 {recvBytes / size:N0} 包(差 {sentBytes / size - recvBytes / size:N0})");
return (sentBytes, recvBytes);
}
/// <summary>往返模å¼ï¼šé€åŒ…å‘é€å¹¶è¯»æ»¡å›žæ˜¾ï¼ˆåŒæ¥/å¼‚æ¥æ‹‰å–æˆ–äº‹ä»¶ä¹’ä¹“ï¼‰ï¼Œé‡‡æ ·å•æ¬¡å¾€è¿”耗时(微秒)</summary>
private static List<Double>? RunRoundTrip(List<ISocketClient> conns, Byte[] payload, Int32 size, Int32 seconds, Boolean tolerateLoss, String recvMode)
{
var samples = new List<Double>();
if (recvMode == "event")
{
// 事件乒乓:å‘é€â†’Received 回执→å†å‘é€ï¼ˆé›¶é˜»å¡žçº¿ç¨‹ï¼›ä¸ŽæœåŠ¡å™¨æŽ¥æ”¶æ¨¡åž‹åŒæž„)
var runners = conns.Select(c => new EventRoundTripClient(c, size, payload, tolerateLoss)).ToArray();
foreach (var r in runners) r.Start();
Thread.Sleep(seconds * 1000);
foreach (var r in runners) r.Stop();
Thread.Sleep(500);
foreach (var r in runners)
{
samples.AddRange(r.Snapshot());
if (r.Missed > 0) Console.WriteLine($"è¦å‘Šï¼š{r.Missed} 次回执超时(丢包容å¿å·²é‡å‘)");
}
}
else
{
var runners = conns.Select(c => new RoundTripClient(c, size, tolerateLoss, recvMode)).ToArray();
var tasks = runners.Select(r => r.Start(payload)).ToArray();
Thread.Sleep(seconds * 1000);
foreach (var r in runners) r.Stop();
var done = Task.WaitAll(tasks, TimeSpan.FromSeconds(30));
if (!done)
{
foreach (var c in conns) c.Close("sleep-timeout");
Task.WaitAll(tasks, TimeSpan.FromSeconds(10));
Console.WriteLine("è¦å‘Šï¼šéƒ¨åˆ†å¾€è¿”ä»»åŠ¡æœªåŠæ—¶ç»“æŸ");
}
foreach (var r in runners) samples.AddRange(r.Samples);
}
if (samples.Count == 0)
{
Console.WriteLine("è¦å‘Šï¼šæœªé‡‡é›†åˆ°å¾€è¿”æ ·æœ¬");
return null;
}
samples.Sort();
var avg = samples.Average();
Console.WriteLine($"å¾€è¿”æ ·æœ¬ : {samples.Count:N0} æ¬¡ï¼ˆçª—å£ {seconds} s)");
Console.WriteLine($"延迟(往返): P50={Percentile(samples, 0.50):F1}µs P95={Percentile(samples, 0.95):F1}µs P99={Percentile(samples, 0.99):F1}µs");
Console.WriteLine($" å¹³å‡={avg:F1}µs 最å°={samples[0]:F1}µs 最大={samples[^1]:F1}µs");
return samples;
}
/// <summary>å–åˆ†ä½æ•°ï¼ˆå‡åºæ ·æœ¬ï¼‰</summary>
private static Double Percentile(List<Double> sorted, Double p)
{
var idx = (Int32)Math.Min(sorted.Count - 1, Math.Max(0, Math.Round((sorted.Count - 1) * p)));
return sorted[idx];
}
/// <summary>å–ä¸ä½æ•°</summary>
private static Double Median(List<Double> values)
{
var sorted = new List<Double>(values);
sorted.Sort();
return Percentile(sorted, 0.50);
}
/// <summary>拉å–一次数æ®é•¿åº¦ï¼ˆåŒæ¥ Receive æˆ–å¼‚æ¥ ReceiveAsyncï¼›UDP 带超时容å¿ä¸¢åŒ…)</summary>
private static async Task<Int32> PullOnce(ISocketClient conn, Boolean useAsync, Boolean tolerateLoss)
{
if (useAsync)
{
if (tolerateLoss)
{
using var cts = new CancellationTokenSource(300);
using var pk = await conn.ReceiveAsync(cts.Token);
return pk?.Length ?? 0;
}
else
{
using var pk = await conn.ReceiveAsync();
return pk?.Length ?? 0;
}
}
else
{
using var pk = conn.Receive();
return pk?.Length ?? 0;
}
}
private static String? GetArg(String[] args, String name)
{
for (var i = 0; i < args.Length - 1; i++)
{
if (args[i] == name) return args[i + 1];
}
return null;
}
private static Int32 GetInt(String[] args, String name, Int32 defaultValue)
{
var text = GetArg(args, name);
return text != null && Int32.TryParse(text, out var value) ? value : defaultValue;
}
/// <summary>æµæ°´çº¿å®¢æˆ·ç«¯ï¼šä¸€ä¸ªå‘é€ä»»åŠ¡ + æŽ¥æ”¶ï¼ˆåŒæ¥/å¼‚æ¥æ‹‰å–循环或事件回调),接收读满å‘逿€»é‡ï¼ˆUDP 容å¿ä¸¢åŒ…)</summary>
private sealed class PipelineClient
{
private readonly ISocketClient _conn;
private readonly Int32 _size;
private readonly Boolean _tolerateLoss;
private readonly Boolean _eventMode;
private readonly Boolean _useAsync;
private volatile Boolean _stopped;
private Int64 _eventBytes;
public Int64 SentPackets;
public Int64 ReceivedBytes;
public PipelineClient(ISocketClient conn, Int32 size, Boolean tolerateLoss, String recvMode)
{
_conn = conn;
_size = size;
_tolerateLoss = tolerateLoss;
_eventMode = recvMode == "event";
_useAsync = recvMode == "asyncpull";
// 事件模å¼ï¼šè®¢é˜…接收环推é€ï¼Œå—节在回调ä¸ç´¯è®¡ï¼ˆä¸ŽæœåС噍åŒä¸€æŽ¥æ”¶æ¨¡åž‹ï¼‰
if (_eventMode) _conn.Received += OnReceived;
}
/// <summary>äº‹ä»¶æ¨¡å¼æŽ¥æ”¶å›žè°ƒï¼ˆæŽ¥æ”¶çŽ¯çº¿ç¨‹ï¼‰</summary>
private void OnReceived(Object? sender, ReceivedEventArgs e)
{
var pk = e.Packet;
if (pk != null && pk.Length > 0) Interlocked.Add(ref _eventBytes, pk.Length);
}
public Task StartSend(Byte[] payload)
=> Task.Run(() =>
{
var n = 0L;
var batch = 0L;
try
{
var errors = 0;
while (!_stopped)
{
// å‘é€å¤±è´¥ï¼ˆå¦‚ UDP 包超é™è¿”回 -1)ä¸è®¡å…¥å‘é€é‡ï¼Œè¿žç»å¤±è´¥åˆ™æå‰é€€å‡ºé˜²æ»è½¬
if (_conn.Send(payload) <= 0)
{
if (++errors > 1000) break;
continue;
}
errors = 0;
n++;
batch += payload.Length;
// æ¯ 256 包汇入一次全局å‘é€è®¡æ•°ï¼Œå…¼é¡¾å®žæ—¶çª—å£ç»Ÿè®¡ä¸Žä½ŽåŽŸå开销
if ((n & 0xFF) == 0)
{
Interlocked.Add(ref _sentBytes, batch);
batch = 0;
}
}
}
finally
{
Interlocked.Add(ref _sentBytes, batch);
SentPackets = n;
}
});
public Task StartReceive()
{
// äº‹ä»¶æ¨¡å¼æ— 接收任务:数æ®ç”±æŽ¥æ”¶çŽ¯å›žè°ƒæŽ¨é€
if (_eventMode) return Task.CompletedTask;
return Task.Run(async () =>
{
var bytes = 0L;
// å‘逿œªç»“æŸå‰æŒç»è¯»å–ï¼›å‘é€ç»“æŸåŽè¯»æ»¡å‘逿€»é‡ï¼ˆæŽ’水)
while (SentPackets == 0 || bytes < SentPackets * _size)
{
try
{
var n = await PullOnce(_conn, _useAsync, _tolerateLoss);
if (n <= 0) break;
bytes += n;
}
catch (SocketException ex) when (_tolerateLoss && ex.SocketErrorCode == SocketError.TimedOut)
{
// UDP æ— æµæŽ§å¯èƒ½ä¸¢åŒ…:å‘é€å·²åœæ¢ä¸”接收超时,视为排水结æŸ
if (SentPackets > 0) break;
}
catch (OperationCanceledException) when (_tolerateLoss)
{
if (SentPackets > 0) break;
}
}
ReceivedBytes = bytes;
});
}
/// <summary>当å‰å·²æŽ¥æ”¶å—节(事件模å¼è¯»å›žè°ƒè®¡æ•°ï¼‰</summary>
public Int64 TotalReceived => _eventMode ? Interlocked.Read(ref _eventBytes) : ReceivedBytes;
public void Stop() => _stopped = true;
}
/// <summary>å¾€è¿”å®¢æˆ·ç«¯ï¼ˆæ‹‰å–æ¨¡å¼ï¼‰ï¼šä¸²è¡Œé€åŒ…å¾€è¿”ï¼Œé‡‡æ ·å•æ¬¡è€—时(微秒);UDP 丢包时跳过该次继ç»</summary>
private sealed class RoundTripClient
{
private readonly ISocketClient _conn;
private readonly Int32 _size;
private readonly Boolean _tolerateLoss;
private readonly Boolean _useAsync;
private volatile Boolean _stopped;
public List<Double> Samples { get; } = [];
public RoundTripClient(ISocketClient conn, Int32 size, Boolean tolerateLoss, String recvMode)
{
_conn = conn;
_size = size;
_tolerateLoss = tolerateLoss;
_useAsync = recvMode == "asyncpull";
}
public Task Start(Byte[] payload)
=> Task.Run(async () =>
{
while (!_stopped)
{
try
{
var t0 = Stopwatch.GetTimestamp();
_conn.Send(payload);
var need = _size;
while (need > 0)
{
var n = await PullOnce(_conn, _useAsync, _tolerateLoss);
if (n <= 0) throw new TimeoutException("å›žæ˜¾ä¸æ–");
need -= n;
}
var us = (Stopwatch.GetTimestamp() - t0) * 1_000_000.0 / Stopwatch.Frequency;
Samples.Add(us);
}
catch (SocketException ex) when (_tolerateLoss && ex.SocketErrorCode == SocketError.TimedOut)
{
// UDP 丢包:跳过该次往返,继ç»ä¸‹ä¸€è½®ï¼ˆå›žæ˜¾é”™ä½ç”±åè®®æœ¬è´¨å†³å®šï¼Œæ ·æœ¬ä»è¿‘似往返时延)
}
catch (OperationCanceledException) when (_tolerateLoss)
{
// å¼‚æ¥æ‹‰å–è¶…æ—¶ç‰ä»·äºŽæŽ¥æ”¶è¶…æ—¶
}
}
});
public void Stop() => _stopped = true;
}
/// <summary>事件模å¼å¾€è¿”客户端:å‘é€â†’Received 回执→å†å‘é€ï¼ˆä¹’乓链),零阻塞线程;UDP 用看门狗é‡å‘容å¿ä¸¢åŒ…</summary>
private sealed class EventRoundTripClient
{
private readonly ISocketClient _conn;
private readonly Int32 _size;
private readonly Byte[] _payload;
private volatile Boolean _stopped;
private Int64 _sentAt;
private Int32 _need;
private readonly Object _lock = new();
public List<Double> Samples { get; } = [];
public Int64 Missed;
public EventRoundTripClient(ISocketClient conn, Int32 size, Byte[] payload, Boolean tolerateLoss)
{
_conn = conn;
_size = size;
_payload = payload;
_conn.Received += OnReceived;
// UDP 丢包容å¿ï¼šé™é»˜ 1 秒视为回执丢失,记录并é‡å‘ï¼ˆä¿æŒä¹’乓链推进)
if (tolerateLoss) _ = Task.Run(Watchdog);
}
public void Start() => SendNext();
private void SendNext()
{
_need = _size;
Volatile.Write(ref _sentAt, Stopwatch.GetTimestamp());
_conn.Send(_payload);
}
/// <summary>回执回调(接收环线程):累计完整回显åŽè®°å½•æ ·æœ¬å¹¶ç«‹å³é“¾å…¥ä¸‹ä¸€æ¬¡å‘é€</summary>
private void OnReceived(Object? sender, ReceivedEventArgs e)
{
var pk = e.Packet;
if (pk == null || pk.Length <= 0) return;
if (_need - pk.Length > 0)
{
_need -= pk.Length;
return;
}
var us = (Stopwatch.GetTimestamp() - Volatile.Read(ref _sentAt)) * 1_000_000.0 / Stopwatch.Frequency;
lock (_lock) Samples.Add(us);
if (_stopped) return;
SendNext();
}
private void Watchdog()
{
while (!_stopped)
{
Thread.Sleep(100);
var sentAt = Volatile.Read(ref _sentAt);
if (sentAt == 0 || Stopwatch.GetTimestamp() - sentAt < Stopwatch.Frequency) continue;
Missed++;
SendNext();
}
}
public void Stop() => _stopped = true;
public List<Double> Snapshot()
{
lock (_lock) return [.. Samples];
}
}
/// <summary>把汇总å—å…¸åºåˆ—化为 JSON 文本</summary>
/// <remarks>æ‰‹å†™è€Œéž System.Text.Json åå°„å¼åºåˆ—化:åŽè€…在 NativeAOT 下默认被ç¦ç”¨ï¼Œ
/// 会让 AOT äº§ç‰©è¾“å‡ºæ±‡æ€»æ—¶ç›´æŽ¥æŠ›å¼‚å¸¸ï¼ˆå¹¶å¸¦æ¥ IL2026/IL3050 告è¦ï¼‰ã€‚键与å–值形æ€å›ºå®šï¼Œæ‰‹å†™è¶³å¤Ÿã€‚</remarks>
/// <param name="dic">键值对</param>
/// <returns>JSON 文本</returns>
private static String ToJson(Dictionary<String, Object?> dic)
{
var sb = Pool.StringBuilder.Get();
sb.Append('{');
var first = true;
foreach (var kv in dic)
{
if (!first) sb.Append(',');
first = false;
sb.Append('"').Append(kv.Key).Append("\":");
switch (kv.Value)
{
case null: sb.Append("null"); break;
case Boolean v: sb.Append(v ? "true" : "false"); break;
case String v: sb.Append('"').Append(v).Append('"'); break;
case IFormattable v: sb.Append(v.ToString(null, CultureInfo.InvariantCulture)); break;
default: sb.Append('"').Append(kv.Value).Append('"'); break;
}
}
sb.Append('}');
return sb.Return(true);
}
}
|