refactor: 消息交付接缝从 SessionBase 下沉 TcpSession 消息泵相关的交付入口(OnMessage/OnMessageAsync/TryMatchResponse)基类自身从不调用,唯一调用方是 TcpSession 的消息泵,随泵下沉到 TcpSession;SessionBase 的匹配器改为 internal 属性 Matcher 供流式会话在帧定界交付时调用,基类内部引用同步收口。 顺带让 Receive/ReceiveAsync 的拉取互斥判据改用 IsReceiving,消除该属性"自身不用"的悬空状态。石头 authored at 2026-10-11 19:51:39
diff --git "a/Doc/\346\266\210\346\201\257\345\215\217\350\256\256\346\240\210.md" "b/Doc/\346\266\210\346\201\257\345\215\217\350\256\256\346\240\210.md"
index 27ba865..9ccd708 100644
--- "a/Doc/\346\266\210\346\201\257\345\215\217\350\256\256\346\240\210.md"
+++ "b/Doc/\346\266\210\346\201\257\345\215\217\350\256\256\346\240\210.md"
@@ -99,7 +99,7 @@ server.Start();
- `Protocol` 须在打开/下发之前设置(`NetServer` 在 `Start` 前设置,随监听服务下发到会话,Unix 域套接字同样透传)。
- 仅流式会话(TCP/SSL/UDS)启动消息泵;`AutoReceive=false` 时不启动。
- 配置透传:`NetServer.Protocol/MaxCache` → `TcpServer` → 会话;`NetClient.Protocol/MatchTimeout` 同理。
-- 消息泵与帧定界配置(`MaxCache`/`MaxFrameSize`/`RequireFullFrame`)属流式会话(`TcpSession`);`SessionBase` 只保留收发原语、协议与匹配队列,经 `StartReceiveFlow`/`StopReceiveFlow`/`RawReceiveTaken` 三个钩子接入消息泵——数据报会话不覆盖即无动作。
+- 消息泵与帧定界配置(`MaxCache`/`MaxFrameSize`/`RequireFullFrame`)与消息交付入口(`OnMessage`/`OnMessageAsync`/`TryMatchResponse`)属流式会话(`TcpSession`);`SessionBase` 只保留收发原语、协议与匹配队列,经 `StartReceiveFlow`/`StopReceiveFlow`/`RawReceiveTaken` 三个钩子接入消息泵——数据报会话不覆盖即无动作。
### 5.2 接收与交付契约
@@ -121,10 +121,10 @@ server.Received += (s, e) =>
```
- **交付契约**:处理器返回后消息收尾(未读体丢弃对齐帧尾、消息释放);`e.Packet`(内存体视图)与消息同生命周期。
-- 同连接消息**串行**处理;需要异步消费流式体时重写会话的 `OnMessageAsync`(消息与数据窗口在其返回前保持有效)。
+- 同连接消息**串行**处理;需要异步消费流式体时重写流式会话(`TcpSession`)的 `OnMessageAsync`(消息与数据窗口在其返回前保持有效)。
- 同步处理器内可同步等待读满:内存体立即完成;流式体的等待属串行语义,正确性优先。
-> **处理入口选择**:`OnReceive(ReceivedEventArgs)` 是事件出口(裸数据与协议消息共用,业务订阅 `Received` 事件);`OnMessage(IMessage)` 为同步处理入口(默认实现构造事件参数并进事件链,零状态机开销);`OnMessageAsync(IMessage)` 为异步入口(默认转调 OnMessage,需要 `await` 流式体或异步业务时重写)。**重写其一即可,不会双触发**。
+> **处理入口选择**(流式会话):`OnReceive(ReceivedEventArgs)` 是事件出口(裸数据与协议消息共用,业务订阅 `Received` 事件);`OnMessage(IMessage)` 为同步处理入口(默认实现构造事件参数并进事件链,零状态机开销);`OnMessageAsync(IMessage)` 为异步入口(默认转调 OnMessage,需要 `await` 流式体或异步业务时重写)。**重写其一即可,不会双触发**。
> **会话处理器(可选)**:服务器可重载 `NetServer.CreateHandler` 为会话挂 `INetHandler`,数据**先经处理器再进入事件**;协议模式下处理器从事件参数转型取消息(`data is ReceivedEventArgs e → e.Message`),返回后消息同样收尾。适合"有状态协议预处理"(如 HTTP 会话解析);普通业务直接用 `Received` 事件即可。
diff --git "a/Doc/\347\275\221\347\273\234\347\274\223\345\206\262\346\211\200\346\234\211\346\235\203\346\236\266\346\236\204.md" "b/Doc/\347\275\221\347\273\234\347\274\223\345\206\262\346\211\200\346\234\211\346\235\203\346\236\266\346\236\204.md"
index 031cf21..8df241f 100644
--- "a/Doc/\347\275\221\347\273\234\347\274\223\345\206\262\346\211\200\346\234\211\346\235\203\346\236\266\346\236\204.md"
+++ "b/Doc/\347\275\221\347\273\234\347\274\223\345\206\262\346\211\200\346\234\211\346\235\203\346\236\266\346\236\204.md"
@@ -49,7 +49,7 @@ flowchart TD
| # | 场景 | 代码位置 | 频率 |
|---|---|---|---|
-| ① | **RPC 响应交付**:命中 await 等待方,响应消息整体交付(等待方 `Dispose` 消息即归还体) | `SessionBase.TryMatchResponse` → `IMatchQueue.Match` 交付等待方任务 | 常见 |
+| ① | **RPC 响应交付**:命中 await 等待方,响应消息整体交付(等待方 `Dispose` 消息即归还体) | `TcpSession.TryMatchResponse`(消息泵交付阶段)→ `IMatchQueue.Match` 交付等待方任务 | 常见 |
| ② | **半包残片保留**:数据不完整时残片留在管道未消费窗口跨轮累积(大帧跨轮必经) | `Pipe` 读取窗口(投递时已转移给管道的共享切片) | 大帧下常见 |
| ③ | **事件处理器主动带出**:事件内把当前帧带出本轮(跨线程/跨 await) | 业务代码:事件内 `e.Packet.Slice(...)`(裸数据模式);协议模式经 `e.Message.Body` 读取或先物化 | 罕见 |
| ✗ | **意外逃逸**:不经 `Slice`,私自留存轮句柄或借用视图 | 任意代码(违规) | 必须评审拦截 |
diff --git a/NewLife.Core/Net/SessionBase.cs b/NewLife.Core/Net/SessionBase.cs
index 2b3f6b8..c62ca97 100644
--- a/NewLife.Core/Net/SessionBase.cs
+++ b/NewLife.Core/Net/SessionBase.cs
@@ -109,19 +109,19 @@ public abstract class SessionBase : DisposeBase, ISocketClient, ITransport, ILog
/// <summary>请求-响应匹配队列。协议模式下等待响应时使用,首次等待时自动创建,可注入共享或自定义实现</summary>
public IMatchQueue? MatchQueue
{
- get => _matcher.Queue;
- set => _matcher.Queue = value;
+ get => Matcher.Queue;
+ set => Matcher.Queue = value;
}
/// <summary>请求-响应匹配等待超时(毫秒)。默认30_000</summary>
public Int32 MatchTimeout
{
- get => _matcher.Timeout;
- set => _matcher.Timeout = value;
+ get => Matcher.Timeout;
+ set => Matcher.Timeout = value;
}
- /// <summary>响应匹配器。等待登记、匹配交付与关闭清理与传输形态无关,流式与数据报会话共用</summary>
- private readonly ResponseMatcher _matcher = new();
+ /// <summary>响应匹配器。等待登记、匹配交付与关闭清理与传输形态无关;流式会话的消息泵在帧定界交付时调用</summary>
+ internal ResponseMatcher Matcher { get; } = new();
/// <summary>最大并发处理数。协议模式下消息处理并发度:1=串行(默认,同连接依次处理);大于1=并行派发(兼作并发上限)</summary>
/// <remarks>
@@ -270,7 +270,7 @@ public abstract class SessionBase : DisposeBase, ISocketClient, ITransport, ILog
await StopReceiveFlowAsync().ConfigureAwait(false);
// 取消挂起的请求-响应等待,避免调用方悬挂
- _matcher.Clear();
+ Matcher.Clear();
var rs = await OnCloseAsync(reason ?? (GetType().Name + "Close"), cancellationToken).ConfigureAwait(false);
@@ -439,7 +439,7 @@ public abstract class SessionBase : DisposeBase, ISocketClient, ITransport, ILog
if (!Open() || Client == null) return null;
// 接收环运行时禁止拉取:两条读路径会争抢同一链路,数据被分流且不可预期
- if (_RecvCount > 0) throw new InvalidOperationException(NoPullMessage);
+ if (IsReceiving) throw new InvalidOperationException(NoPullMessage);
return OnDirectReceive();
}
@@ -490,7 +490,7 @@ public abstract class SessionBase : DisposeBase, ISocketClient, ITransport, ILog
if (!Open() || Client == null) return null;
// 接收环运行时禁止拉取:两条读路径会争抢同一链路,数据被分流且不可预期
- if (_RecvCount > 0) throw new InvalidOperationException(NoPullMessage);
+ if (IsReceiving) throw new InvalidOperationException(NoPullMessage);
return await OnDirectReceiveAsync(cancellationToken).ConfigureAwait(false);
}
@@ -938,57 +938,6 @@ public abstract class SessionBase : DisposeBase, ISocketClient, ITransport, ILog
#endregion 接收
- #region 消息交付
-
- /// <summary>尝试把响应消息匹配给等待中的请求(协议模式)。命中则交付等待方,跳过收尾</summary>
- /// <param name="message">收到的消息</param>
- /// <returns>是否已匹配交付</returns>
- /// <remarks>
- /// <para>无匹配队列(未发过等待请求)或无配对协议时快速返回,不产生额外开销;
- /// 消息是否可配对由协议 matcher 判定(如 SRMP 要求应答消息+序列号相等;无方向协议可用恒真 matcher)。</para>
- /// <para>流式负载在事件交付前物化为内存模式:事件链可观察读取,等待方在任意时机异步消费(一次拷贝换正确性)。
- /// 未命中时消息按普通流程收尾,负载随消息归还。</para>
- /// </remarks>
- protected Boolean TryMatchResponse(IMessage message) => _matcher.TryMatch(this, Protocol, message);
-
- /// <summary>收到消息(异步)。协议模式(<see cref="Protocol"/> 非空)下由消息泵逐帧调用</summary>
- /// <param name="message">消息(头部字段就位、体已绑定)</param>
- /// <remarks>
- /// <para>默认触发同步 <see cref="Received"/> 事件链。处理器返回后消息进入收尾:未读体被丢弃对齐帧尾、消息释放。</para>
- /// <para>需要异步读取流式主体的场景,继承会话重写本方法,在 await 期间消息与数据窗口保持有效;同步事件处理器内需要流式数据时请先物化(<see cref="LimitedReader.ReadAllAsync"/>,数据未到齐会等待——串行语义下正确),或物化后交给后台异步链处理。</para>
- /// </remarks>
- protected virtual ValueTask OnMessageAsync(IMessage message)
- {
- OnMessage(message);
-
- return default;
- }
-
- /// <summary>收到消息。协议模式(<see cref="Protocol"/> 非空)下由消息泵逐帧调用</summary>
- /// <param name="message">消息(头部字段就位、体已绑定)</param>
- /// <remarks>
- /// <para>构造接收事件参数并进入 <see cref="OnReceive"/> 事件链;消息体未读部分由消息泵在交付后丢弃对齐下一帧。</para>
- /// <para>事件参数的 <see cref="ReceivedEventArgs.Packet"/> 为消息负载视图(整帧快路径),流式模式下为 null;与消息同生命周期(处理器返回后失效),需要留存请先物化或切出共享句柄。业务请统一经 <see cref="ReceivedEventArgs.Message"/> 读取头部与流式体。</para>
- /// </remarks>
- protected virtual void OnMessage(IMessage message)
- {
- var e = ReceivedEventArgs.Rent();
- try
- {
- e.Local = Local.Address;
- e.Remote = Remote.EndPoint;
- e.Packet = message.Payload;
- e.Message = message;
-
- OnReceive(e);
- }
- finally
- {
- ReceivedEventArgs.Return(e);
- }
- }
- #endregion
-
#region 消息处理
/// <summary>发送消息。经协议(<see cref="Protocol"/>)构建整帧后发送,不等待响应</summary>
@@ -1043,7 +992,7 @@ public abstract class SessionBase : DisposeBase, ISocketClient, ITransport, ILog
if (this is not IStreamSession) throw new NotSupportedException($"会话类型 [{GetType().Name}] 未接入响应匹配交付,不支持请求-响应等待(数据报会话请用 UdpServer/UdpSession)");
var span = Tracer?.NewSpan($"net:{Name}:SendMessageAsync", request);
- var source = _matcher.Register(this, request, span);
+ var source = Matcher.Register(this, request, span);
try
{
diff --git a/NewLife.Core/Net/TcpSession.Pump.cs b/NewLife.Core/Net/TcpSession.Pump.cs
index 3da4211..828b5da 100644
--- a/NewLife.Core/Net/TcpSession.Pump.cs
+++ b/NewLife.Core/Net/TcpSession.Pump.cs
@@ -5,8 +5,9 @@ namespace NewLife.Net;
/// <summary>增强TCP客户端(协议模式的消息泵与帧定界交付)</summary>
/// <remarks>
-/// <para>协议模式下数据经数据管道定界,由消息泵逐帧交付。流式专有配置(最大残留、单帧上限、整帧模式)与并发派发信号量随本类,
-/// 不再放到 <see cref="SessionBase"/>:数据报会话没有粘包与定界问题,这些成员在它那里只会是"未使用"。</para>
+/// <para>协议模式下数据经数据管道定界,由消息泵逐帧交付。流式专有配置(最大残留、单帧上限、整帧模式)、并发派发信号量
+/// 与消息交付入口(<see cref="OnMessageAsync"/> 等)随本类,不再放到 <see cref="SessionBase"/>:
+/// 数据报会话没有粘包与定界问题,这些成员在它那里只会是"未使用"。</para>
/// </remarks>
public partial class TcpSession
{
@@ -350,4 +351,56 @@ public partial class TcpSession
private SemaphoreSlim Concurrency => _concurrency ??= new SemaphoreSlim(MaxConcurrency, MaxConcurrency);
#endregion
+
+ #region 消息交付
+
+ /// <summary>尝试把响应消息匹配给等待中的请求(协议模式)。命中则交付等待方,跳过收尾</summary>
+ /// <param name="message">收到的消息</param>
+ /// <returns>是否已匹配交付</returns>
+ /// <remarks>
+ /// <para>无匹配队列(未发过等待请求)或无配对协议时快速返回,不产生额外开销;
+ /// 消息是否可配对由协议 matcher 判定(如 SRMP 要求应答消息+序列号相等;无方向协议可用恒真 matcher)。</para>
+ /// <para>流式负载在事件交付前物化为内存模式:事件链可观察读取,等待方在任意时机异步消费(一次拷贝换正确性)。
+ /// 未命中时消息按普通流程收尾,负载随消息归还。</para>
+ /// </remarks>
+ private Boolean TryMatchResponse(IMessage message) => Matcher.TryMatch(this, Protocol, message);
+
+ /// <summary>收到消息(异步)。协议模式(<see cref="SessionBase.Protocol"/> 非空)下由消息泵逐帧调用</summary>
+ /// <param name="message">消息(头部字段就位、体已绑定)</param>
+ /// <remarks>
+ /// <para>默认触发同步 <see cref="SessionBase.Received"/> 事件链。处理器返回后消息进入收尾:未读体被丢弃对齐帧尾、消息释放。</para>
+ /// <para>需要异步读取流式主体的场景,继承会话重写本方法,在 await 期间消息与数据窗口保持有效;同步事件处理器内需要流式数据时请先物化(<see cref="LimitedReader.ReadAllAsync"/>,数据未到齐会等待——串行语义下正确),或物化后交给后台异步链处理。</para>
+ /// </remarks>
+ protected virtual ValueTask OnMessageAsync(IMessage message)
+ {
+ OnMessage(message);
+
+ return default;
+ }
+
+ /// <summary>收到消息。协议模式(<see cref="SessionBase.Protocol"/> 非空)下由消息泵逐帧调用</summary>
+ /// <param name="message">消息(头部字段就位、体已绑定)</param>
+ /// <remarks>
+ /// <para>构造接收事件参数并进入 <see cref="SessionBase.OnReceive"/> 事件链;消息体未读部分由消息泵在交付后丢弃对齐下一帧。</para>
+ /// <para>事件参数的 <see cref="ReceivedEventArgs.Packet"/> 为消息负载视图(整帧快路径),流式模式下为 null;与消息同生命周期(处理器返回后失效),需要留存请先物化或切出共享句柄。业务请统一经 <see cref="ReceivedEventArgs.Message"/> 读取头部与流式体。</para>
+ /// </remarks>
+ protected virtual void OnMessage(IMessage message)
+ {
+ var e = ReceivedEventArgs.Rent();
+ try
+ {
+ e.Local = Local.Address;
+ e.Remote = Remote.EndPoint;
+ e.Packet = message.Payload;
+ e.Message = message;
+
+ OnReceive(e);
+ }
+ finally
+ {
+ ReceivedEventArgs.Return(e);
+ }
+ }
+
+ #endregion
}