using System.Buffers;
using System.ComponentModel;
using NewLife;
using NewLife.Data;
using NewLife.Messaging;
using NewLife.Net;
using Xunit;
namespace XUnitTest.Net;
/// <summary>ä¼šè¯æ•°æ®ç®¡é“(IStreamSession.Pipe æµå¼æŽ¥æ”¶ï¼‰æµ‹è¯•</summary>
[Collection("Net")]
public class SessionPipeTests
{
[Fact]
[DisplayName("会è¯ç®¡é“_接收投递_帧泵æµå¼è¯»å–")]
public async Task SessionPipe_Frames()
{
using var server = new NetServer { Port = 0 };
server.Start();
var wait = new ManualResetEventSlim();
Pipe? pipe = null;
server.NewSession += (s, e) =>
{
pipe = (e.Session.Session as IStreamSession)?.Pipe;
wait.Set();
};
using var client = new NetUri($"tcp://127.0.0.1:{server.Port}").CreateRemote();
client.Open();
// 会è¯å»ºç«‹ä¸”管é“就绪åŽå†å‘é€ï¼Œé¿å…é¦–è½®æ•°æ®æ—©äºŽç®¡é“创建
Assert.True(wait.Wait(3_000));
Assert.NotNull(pipe);
// ä¸¤å¸§ï¼ˆæ ‡å‡† SRMP):第二帧分两次å‘é€ï¼ŒéªŒè¯è·¨è½®ç»„帧
var f1 = BuildFrame(Fill(10, 1));
var f2 = BuildFrame(Fill(300, 2));
_ = client.Send(f1);
_ = client.Send(f2[..100]);
await Task.Delay(20);
_ = client.Send(f2[100..]);
var pump = new MessagePump(new SrmpCodec());
var m1 = await pump.ReadAsync(pipe!.Reader).AsTask().WaitAsync(TimeSpan.FromSeconds(5));
Assert.NotNull(m1);
var b1 = await m1!.Body!.ReadAllAsync();
Assert.Equal(Fill(10, 1), b1.AsReadOnlySequence().ToArray());
b1.TryDispose();
m1.Dispose();
var m2 = await pump.ReadAsync(pipe!.Reader).AsTask().WaitAsync(TimeSpan.FromSeconds(5));
Assert.NotNull(m2);
var b2 = await m2!.Body!.ReadAllAsync();
Assert.Equal(Fill(300, 2), b2.AsReadOnlySequence().ToArray());
b2.TryDispose();
m2.Dispose();
}
[Fact]
[DisplayName("会è¯ç®¡é“_å…³é—_通知消费方æµç»“æŸ")]
public async Task SessionPipe_Close_CompletesReader()
{
using var server = new NetServer { Port = 0 };
server.Start();
var wait = new ManualResetEventSlim();
Pipe? pipe = null;
server.NewSession += (s, e) =>
{
pipe = (e.Session.Session as IStreamSession)?.Pipe;
wait.Set();
};
using var client = new NetUri($"tcp://127.0.0.1:{server.Port}").CreateRemote();
client.Open();
Assert.True(wait.Wait(3_000));
Assert.NotNull(pipe);
// 挂起读å–,客户端æ–å¼€åŽæœåŠ¡ç«¯ä¼šè¯å…³é—,读å–应立å³å®Œæˆï¼ˆIsCompleted)
var vt = pipe!.Reader.ReadAsync();
Assert.False(vt.IsCompleted);
client.Close("test");
var rr = await vt.AsTask().WaitAsync(TimeSpan.FromSeconds(5));
Assert.True(rr.IsCompleted);
Assert.True(rr.Buffer.IsEmpty);
}
[Fact]
[DisplayName("会è¯ç®¡é“_背压_è¾¾åˆ°æš‚åœæ°´ä½åŽæ¢å¤")]
public async Task SessionPipe_Backpressure_PauseResume()
{
using var server = new NetServer { Port = 0 };
server.Start();
var wait = new ManualResetEventSlim();
Pipe? pipe = null;
server.NewSession += (s, e) =>
{
pipe = (e.Session.Session as IStreamSession)?.Pipe;
wait.Set();
};
using var client = new NetUri($"tcp://127.0.0.1:{server.Port}").CreateRemote();
client.Open();
Assert.True(wait.Wait(3_000));
Assert.NotNull(pipe);
// 自定义低水ä½ï¼Œé¿å…çœŸå‘ 1MB æ‰èƒ½è§¦å‘
pipe!.PauseThreshold = 32 * 1024;
pipe.ResumeThreshold = 16 * 1024;
var resumed = 0;
pipe.Resumed += (s, e) => Interlocked.Increment(ref resumed);
// å‘é€ç«¯åŽå°æŽ¨é€ 40 帧 × 4KBï¼šæŽ¥æ”¶ç«¯ä¸æ¶ˆè´¹æ—¶ç®¡é“åº”è¾¾åˆ°æš‚åœæ°´ä½å¹¶åœæ¢æŽ¥æ”¶
const Int32 frameCount = 40;
var frames = new List<Byte[]>();
for (var i = 0; i < frameCount; i++)
{
var f = BuildFrame(Fill(4 * 1024, (Byte)i));
frames.Add(f);
}
var sender = Task.Run(() =>
{
foreach (var f in frames) client.Send(f);
});
// ç‰å¾…管é“è¾¾åˆ°æš‚åœæ°´ä½ï¼ˆæŽ¥æ”¶æ–¹æœªæ¶ˆè´¹ï¼Œæ•°æ®æŒç»ç´¯ç§¯ï¼‰
var sw = new System.Diagnostics.Stopwatch();
sw.Start();
while (!pipe.IsPaused && sw.Elapsed < TimeSpan.FromSeconds(10)) await Task.Delay(10);
Assert.True(pipe.IsPaused, "æŽ¥æ”¶æ•°æ®æœªè§¦å‘æš‚åœæ°´ä½");
// æŒç»æ¶ˆè´¹ï¼šé™åˆ°æ¢å¤æ°´ä½ä»¥ä¸‹åº”è§¦å‘ Resumed å¹¶æ¢å¤æŽ¥æ”¶ï¼Œæœ€ç»ˆæ”¶å®Œæ‰€æœ‰å¸§
var pump = new MessagePump(new SrmpCodec());
var received = 0;
var total = 0L;
while (received < frameCount)
{
var m = await pump.ReadAsync(pipe.Reader).AsTask().WaitAsync(TimeSpan.FromSeconds(10));
Assert.NotNull(m);
var body = await m!.Body!.ReadAllAsync();
total += body.Total;
body.TryDispose();
m.Dispose();
received++;
}
await sender.WaitAsync(TimeSpan.FromSeconds(10));
Assert.True(resumed > 0, "消费é™åŽ‹åŽæœªè§¦å‘æ¢å¤äº‹ä»¶");
Assert.Equal((Int64)frameCount * 4 * 1024, total);
Assert.False(pipe.IsPaused);
}
[Fact]
[DisplayName("会è¯ç®¡é“_整帧路径_超大帧凑é½ä¸å—æš‚åœé˜»æŒ¡ï¼ˆè¯»é¥¥é¥¿è®©ä½ï¼‰")]
public async Task SessionPipe_WholeFrameExceedsPauseThreshold_ReaderStarvationReleases()
{
using var server = new NetServer { Port = 0 };
server.Start();
var wait = new ManualResetEventSlim();
Pipe? pipe = null;
server.NewSession += (s, e) =>
{
pipe = (e.Session.Session as TcpSession)?.Pipe;
wait.Set();
};
using var client = new NetUri($"tcp://127.0.0.1:{server.Port}").CreateRemote();
client.Open();
Assert.True(wait.Wait(3_000));
Assert.NotNull(pipe);
// 低水ä½ï¼ˆ64K æš‚åœ / 32K æ¢å¤ï¼‰ï¼›å•帧声明 256Kï¼Œè¿œè¶…æš‚åœæ°´ä½
pipe!.PauseThreshold = 64 * 1024;
pipe.ResumeThreshold = 32 * 1024;
var resumed = 0;
pipe.Resumed += (s, e) => Interlocked.Increment(ref resumed);
var frame = BuildFrame(Fill(256 * 1024, 0x5A));
// åŽå°åˆ†å—推é€ï¼ˆæ¨¡æ‹Ÿç½‘络分片);远超水ä½åŽ send å¯èƒ½åœé¡¿äºŽå†…æ ¸ç¼“å†²ï¼Œå±žé¢„æœŸ
var sender = Task.Run(() =>
{
try
{
for (var i = 0; i < frame.Length; i += 8 * 1024)
{
var count = Math.Min(8 * 1024, frame.Length - i);
client.Send(frame, i, count);
}
}
catch { }
});
// 整帧路径:帧ä¸å®Œæ•´ä¸æ¶ˆè´¹ï¼Œæš‚åœå…ˆè§¦å‘;但读侧挂起ç‰å¾…ï¼ˆé¥¥é¥¿ï¼‰æ—¶è®©ä½æ”¾è¡Œï¼Œå¸§åº”完整到达
var pump = new MessagePump(new SrmpCodec());
var m = await pump.ReadAsync(pipe.Reader).AsTask().WaitAsync(TimeSpan.FromSeconds(15));
Assert.NotNull(m);
var body = await m!.Body!.ReadAllAsync();
Assert.Equal(256 * 1024, body.Total);
Assert.Equal(frame[8..], body.AsReadOnlySequence().ToArray());
body.TryDispose();
m.Dispose();
// 读饥饿让ä½è‡³å°‘å‘生一次:暂åœå·²è§¦å‘,但读者ç‰å¾…的数æ®è¢«æ”¾è¡Œ
Assert.True(resumed > 0, "超大帧ç‰å¾…期间未å‘生读饥饿让ä½");
await Task.WhenAny(sender, Task.Delay(2_000));
}
[Fact]
[DisplayName("会è¯ç®¡é“_背压_é«˜é¢‘æš‚åœæ¢å¤_å¤šè½®æ— ä¸¢å”¤é†’")]
public async Task SessionPipe_Backpressure_HighFrequencyPauseResume()
{
using var server = new NetServer { Port = 0 };
server.Start();
var wait = new ManualResetEventSlim();
Pipe? pipe = null;
server.NewSession += (s, e) =>
{
pipe = (e.Session.Session as TcpSession)?.Pipe;
wait.Set();
};
using var client = new NetUri($"tcp://127.0.0.1:{server.Port}").CreateRemote();
client.Open();
Assert.True(wait.Wait(3_000));
Assert.NotNull(pipe);
// å°æ°´ä½ + å¿«æ¶ˆè´¹ï¼šåˆ¶é€ é«˜é¢‘æš‚åœ/æ¢å¤å¾ªçŽ¯ï¼Œè¦†ç›–â€œæš‚å˜æ™šäºŽæ¶ˆè´¹æ¢å¤â€çš„丢唤醒窗å£
pipe!.PauseThreshold = 16 * 1024;
pipe.ResumeThreshold = 8 * 1024;
var pump = new MessagePump(new SrmpCodec());
// 多轮å‘é€ï¼šè·¨è½®å¤çŽ°â€œå‰è½®æ£å¸¸ã€åŽè½®æŽ¥æ”¶åœæ‘†â€çš„场景(对é½åŸºå‡†è§„æ¨¡ï¼Œåˆ¶é€ é«˜é¢‘æš‚åœ/æ¢å¤ï¼‰
for (var round = 0; round < 2; round++)
{
const Int32 frameCount = 1024;
var frame = BuildFrame(Fill(16 * 1024, (Byte)round));
var sender = Task.Run(() =>
{
try
{
for (var i = 0; i < frameCount; i++) client.Send(frame);
}
catch { }
});
var received = 0;
try
{
while (received < frameCount)
{
var m = await pump.ReadAsync(pipe.Reader).AsTask().WaitAsync(TimeSpan.FromSeconds(30));
Assert.NotNull(m);
var body = await m!.Body!.ReadAllAsync();
body.TryDispose();
m.Dispose();
received++;
}
}
catch (TimeoutException)
{
throw new TimeoutException($"第 {round} è½®æŽ¥æ”¶åœæ‘†ï¼šå·²æ”¶ {received}/{frameCount},暂åœ={pipe.IsPaused},未消费={pipe.UnconsumedLength}");
}
try
{
await sender.WaitAsync(TimeSpan.FromSeconds(30));
}
catch (TimeoutException)
{
throw new TimeoutException($"第 {round} è½®å‘逿œªå®Œæˆï¼šæš‚åœ={pipe.IsPaused},未消费={pipe.UnconsumedLength}");
}
}
client.Close("test");
}
#region 工具
private static Byte[] Fill(Int32 count, Byte value)
{
var buf = new Byte[count];
for (var i = 0; i < count; i++) buf[i] = value;
return buf;
}
/// <summary>æž„é€ æ ‡å‡†æ¶ˆæ¯å¸§ï¼ˆ4/8å—节头 + 负载)</summary>
private static Byte[] BuildFrame(Byte[] payload)
{
var headerSize = payload.Length < 0xFFFF ? 4 : 8;
var buf = new Byte[headerSize + payload.Length];
buf[0] = 0x01;
buf[1] = 0x02;
if (headerSize == 4)
{
buf[2] = (Byte)(payload.Length & 0xFF);
buf[3] = (Byte)(payload.Length >> 8);
}
else
{
buf[2] = 0xFF;
buf[3] = 0xFF;
buf[4] = (Byte)(payload.Length & 0xFF);
buf[5] = (Byte)((payload.Length >> 8) & 0xFF);
buf[6] = (Byte)((payload.Length >> 16) & 0xFF);
buf[7] = (Byte)((payload.Length >> 24) & 0xFF);
}
payload.CopyTo(buf, headerSize);
return buf;
}
#endregion
}
|