自定义 Provider
本文以 Unity WebSocket 为例,说明如何为 KBEngine Nex C# SDK 创建一个可用于生产环境的网络层 Provider。阅读前建议先了解《网络层 Provider》中的职责边界、网络模式和接入方式。
完整参考实现:
示例使用 NativeWebSocket 作为底层传输,因为它同时支持 Unity Editor、Standalone 和 WebGL。Provider 接口本身不依赖 NativeWebSocket;接入微信小游戏、抖音小游戏、客户端引擎自带 Socket 或其他网络库时,只需替换底层网络调用,并保留本文的生命周期、队列和内存所有权设计。
工程结构
建议把 Provider 和第三方传输放在独立目录中:
UnityWebSocketProvider/
├── UnityWebSocketNetworkProvider.cs # Provider 工厂与单连接实现
├── WebSocket.cs # NativeWebSocket C# 实现
├── WebSocket.jslib # WebGL 构建期桥接
├── NativeWebSocket.LICENSE.txt # 第三方许可证
└── THIRD_PARTY_NOTICES.md # 第三方来源和维护说明Provider 分为两层:
UnityWebSocketProviderFactory
保存不可变配置
每次 Create(context) 创建一个新实例
|
v
UnityWebSocketNetworkProvider
一个实例只对应一条 Loginapp 或 Baseapp 连接
管理状态、事件队列、发送队列和关闭流程
|
v
NativeWebSocket.WebSocket
执行实际连接、发送和关闭不要让工厂持有 WebSocket 连接。Loginapp、Baseapp 和 Relogin 生命周期互相独立,复用实例会让旧回调、发送队列和关闭状态污染新 Session。
第一步:理解 SDK 接口
Provider 需要实现以下接口:
using System;
using System.Collections.Generic;
namespace KBEngine
{
public interface INetworkProviderFactory
{
INetworkProvider Create(NetworkProviderContext context);
}
public interface INetworkProvider : IDisposable
{
NetworkProviderState State { get; }
void Connect(INetworkProviderListener listener);
NetworkSendResult Send(IReadOnlyList<ArraySegment<byte>> packets);
void Process();
void Close();
}
public interface INetworkProviderListener
{
void OnConnected();
void OnDataReceived(byte[] buffer, int offset, int count);
void OnClosed(NetworkCloseInfo closeInfo);
}
}调用关系如下:
| 方法 | 调用方 | Provider 的职责 |
|---|---|---|
Factory.Create() | SDK | 根据 Context 创建全新连接实例 |
Connect() | SDK | 保存 Listener,启动一次异步连接 |
Send() | SDK | 整批接收或拒绝字节,并返回发送结果 |
Process() | SDK Tick | 驱动底层消息队列,派发已排队事件 |
Close() / Dispose() | SDK | 幂等停止连接并释放全部资源 |
Listener.On*() | Provider | 只能在 Process() 内同步调用 |
Provider 交付的是有序字节流,不是 KBEngine 消息。它不解析消息 ID、消息长度、Entity 或加密内容。TCP 读取块和 WebSocket Message 都不等于 KBEngine 协议消息,拆包和粘包由 SDK 处理。
第二步:实现工厂
工厂保存所有连接共享的配置,并为每个 Context 创建独立 Provider:
namespace KBEngine.UnityWebSocket
{
using System;
using System.Collections.Generic;
public sealed class UnityWebSocketProviderFactory : INetworkProviderFactory
{
private readonly bool _secure;
private readonly string _path;
private readonly Dictionary<string, string> _domainMapping;
private readonly Dictionary<int, int> _portMapping;
public UnityWebSocketProviderFactory(
bool secure,
string path,
IReadOnlyDictionary<string, string> domainMapping,
IReadOnlyDictionary<int, int> portMapping)
{
_secure = secure;
_path = string.IsNullOrWhiteSpace(path) ? "/" : path.Trim();
_domainMapping = SnapshotDomainMapping(domainMapping);
_portMapping = SnapshotPortMapping(portMapping);
}
public INetworkProvider Create(NetworkProviderContext context)
{
if (context == null)
throw new ArgumentNullException(nameof(context));
// 每个 Session 使用独立实例,避免连接状态跨 Session 泄漏。
// Each Session gets an isolated instance to prevent state leakage.
return new UnityWebSocketNetworkProvider(
context, _secure, _path, _domainMapping, _portMapping);
}
private static Dictionary<string, string> SnapshotDomainMapping(
IReadOnlyDictionary<string, string> source)
{
var result = new Dictionary<string, string>(StringComparer.OrdinalIgnoreCase);
if (source == null)
return result;
foreach (KeyValuePair<string, string> item in source)
{
if (string.IsNullOrWhiteSpace(item.Key) ||
string.IsNullOrWhiteSpace(item.Value))
{
throw new ArgumentException(
"Domain mappings cannot contain empty hosts.", nameof(source));
}
result.Add(item.Key.Trim(), item.Value.Trim());
}
return result;
}
private static Dictionary<int, int> SnapshotPortMapping(
IReadOnlyDictionary<int, int> source)
{
var result = new Dictionary<int, int>();
if (source == null)
return result;
foreach (KeyValuePair<int, int> item in source)
{
if (item.Key <= 0 || item.Key > ushort.MaxValue ||
item.Value <= 0 || item.Value > ushort.MaxValue)
{
throw new ArgumentException(
"Port mappings must contain ports between 1 and 65535.",
nameof(source));
}
result.Add(item.Key, item.Value);
}
return result;
}
}
}工厂构造时复制配置,而不是保存调用者的字典引用。这样外部修改不会影响已经创建的连接,也避免并发读取字典时被其他线程修改。非法映射应尽早抛出,不能等到握手失败后再用模糊日志排查。
第三步:定义状态和队列
Provider 需要一个生命周期锁和一个事件队列锁:
namespace KBEngine.UnityWebSocket
{
using System;
using System.Collections.Generic;
using NativeWebSocket;
public sealed class UnityWebSocketNetworkProvider : INetworkProvider
{
private readonly object _lifecycleLock = new object();
private readonly object _eventLock = new object();
private readonly Queue<ProviderEvent> _events = new Queue<ProviderEvent>();
private readonly Queue<byte[]> _sendQueue = new Queue<byte[]>();
private readonly NetworkProviderContext _context;
private readonly string _url;
private INetworkProviderListener _listener;
private WebSocket _socket;
private volatile NetworkProviderState _state = NetworkProviderState.Created;
private int _queuedSendBytes;
private bool _sendWorkerRunning;
private bool _disposed;
public UnityWebSocketNetworkProvider(
NetworkProviderContext context,
bool secure,
string path,
IReadOnlyDictionary<string, string> domainMapping,
IReadOnlyDictionary<int, int> portMapping)
{
_context = context ?? throw new ArgumentNullException(nameof(context));
if (context.SendQueueSize <= 0)
throw new ArgumentOutOfRangeException(nameof(context));
_url = BuildUrl(context.Endpoint, secure, path, domainMapping, portMapping);
}
public NetworkProviderState State => _state;
}
}状态机为:
Created -> Connecting -> Connected -> Closing -> Closed
| |
+-----------+------> Failedvolatile 只保证读取 State 时能看到最新值,不能代替锁。检查状态、替换 Socket、更新队列必须在同一个 _lifecycleLock 临界区完成。
第四步:从 Endpoint 构造 URL
Provider 必须以 SDK 提供的 context.Endpoint 为连接源:
private static string BuildUrl(
NetworkEndpoint endpoint,
bool secure,
string path,
IReadOnlyDictionary<string, string> domainMapping,
IReadOnlyDictionary<int, int> portMapping)
{
if (endpoint == null)
throw new ArgumentNullException(nameof(endpoint));
string host = endpoint.Host;
int port = endpoint.TcpPort;
if (domainMapping != null && domainMapping.TryGetValue(host, out string mappedHost))
host = mappedHost;
if (portMapping != null && portMapping.TryGetValue(port, out int mappedPort))
port = mappedPort;
if (port <= 0 || port > ushort.MaxValue)
throw new InvalidOperationException("The resolved port is invalid.");
string normalizedPath = string.IsNullOrWhiteSpace(path) ? "/" : path.Trim();
if (normalizedPath.IndexOf('#') >= 0)
throw new InvalidOperationException("WebSocket paths cannot contain a fragment.");
string query = string.Empty;
int queryIndex = normalizedPath.IndexOf('?');
if (queryIndex >= 0)
{
query = normalizedPath.Substring(queryIndex + 1);
normalizedPath = normalizedPath.Substring(0, queryIndex);
}
if (!normalizedPath.StartsWith("/", StringComparison.Ordinal))
normalizedPath = "/" + normalizedPath;
var builder = new UriBuilder(secure ? "wss" : "ws", host, port)
{
Path = normalizedPath,
Query = query,
};
return builder.Uri.AbsoluteUri;
}使用 UriBuilder 可以正确处理 IPv6、路径和查询参数。WebSocket 使用 Endpoint.TcpPort,因为 KBEngine 服务端在外部 TCP Channel 上识别 WebSocket Upgrade。Provider 只能转换公开入口,不能写死 Baseapp 地址。
第五步:建立连接
public void Connect(INetworkProviderListener listener)
{
if (listener == null)
throw new ArgumentNullException(nameof(listener));
WebSocket socket;
lock (_lifecycleLock)
{
if (_disposed)
throw new ObjectDisposedException(nameof(UnityWebSocketNetworkProvider));
if (_state != NetworkProviderState.Created)
throw new InvalidOperationException("A provider can connect only once.");
socket = new WebSocket(_url);
socket.OnOpen += OnSocketOpen;
socket.OnMessage += OnSocketMessage;
socket.OnError += OnSocketError;
socket.OnClose += OnSocketClose;
_listener = listener;
_socket = socket;
_state = NetworkProviderState.Connecting;
}
_ = ConnectAsync(socket);
}
private async System.Threading.Tasks.Task ConnectAsync(WebSocket socket)
{
try
{
await socket.Connect();
}
catch (Exception exception)
{
Fail("WebSocket connect failed.", exception);
}
}不要使用 async void,否则连接异常不能统一进入 Provider 的失败流程。Connect() 也不能同步等待握手,否则可能阻塞 Unity 主线程或 WebGL 运行时。
第六步:底层回调只负责排队
平台回调可能来自任意线程,不能直接调用 SDK Listener:
private void OnSocketOpen()
{
lock (_lifecycleLock)
{
if (_disposed || _state != NetworkProviderState.Connecting)
return;
_state = NetworkProviderState.Connected;
EnqueueEvent(ProviderEvent.Connected());
}
}
private void OnSocketMessage(byte[] data)
{
if (data == null || data.Length == 0)
return;
lock (_lifecycleLock)
{
if (_disposed || _state != NetworkProviderState.Connected)
return;
EnqueueEvent(ProviderEvent.Data(data));
}
}
private void OnSocketError(string message)
{
string detail = string.IsNullOrWhiteSpace(message)
? "Unknown WebSocket error."
: message;
Fail("WebSocket transport error: " + detail, new InvalidOperationException(detail));
}
private void OnSocketClose(WebSocketCloseCode code)
{
Fail("WebSocket closed with code " + (int)code + " (" + code + ").", null);
}
private void EnqueueEvent(ProviderEvent item)
{
lock (_eventLock)
_events.Enqueue(item);
}事件载体如下:
private enum ProviderEventType { Connected, Data, Closed }
private sealed class ProviderEvent
{
private ProviderEvent(
ProviderEventType type, byte[] data, NetworkCloseInfo closeInfo)
{
Type = type;
DataBuffer = data;
CloseInfo = closeInfo;
}
internal ProviderEventType Type { get; }
internal byte[] DataBuffer { get; }
internal NetworkCloseInfo CloseInfo { get; }
internal static ProviderEvent Connected() =>
new ProviderEvent(ProviderEventType.Connected, null, null);
internal static ProviderEvent Data(byte[] data) =>
new ProviderEvent(ProviderEventType.Data, data, null);
internal static ProviderEvent Closed(NetworkCloseInfo closeInfo) =>
new ProviderEvent(ProviderEventType.Closed, null, closeInfo);
}如果底层网络库会复用接收缓冲,必须先复制再排队。不能保存只在平台回调期间有效的指针、ArrayBuffer 视图或临时 Buffer。
第七步:在 Process 中派发
public void Process()
{
WebSocket socket = _socket;
#if !UNITY_WEBGL || UNITY_EDITOR
if (socket != null &&
(_state == NetworkProviderState.Connecting ||
_state == NetworkProviderState.Connected))
{
try
{
socket.DispatchMessageQueue();
}
catch (Exception exception)
{
Fail("WebSocket message dispatch failed.", exception);
}
}
#endif
DrainEvents();
}
private void DrainEvents()
{
while (true)
{
ProviderEvent item;
lock (_eventLock)
{
if (_events.Count == 0)
return;
item = _events.Dequeue();
}
// Listener 必须在锁外调用,避免 SDK 回调重入造成死锁。
// Invoke listeners outside locks to avoid deadlocks during SDK re-entry.
switch (item.Type)
{
case ProviderEventType.Connected:
_listener?.OnConnected();
break;
case ProviderEventType.Data:
_listener?.OnDataReceived(
item.DataBuffer, 0, item.DataBuffer.Length);
break;
case ProviderEventType.Closed:
_listener?.OnClosed(item.CloseInfo);
break;
}
}
}Process() 是平台事件进入 SDK 的唯一入口。NativeWebSocket 在 Editor/Standalone 需要 DispatchMessageQueue(),WebGL 使用不同实现,因此要保留条件编译。若单 Tick 消息量可能很大,应增加数量或耗时预算,但必须保持 FIFO,并把未处理事件留到下一 Tick。
第八步:实现发送和背压
SDK 传入的 ArraySegment<byte> 可能指向对象池。返回 Accepted 前必须复制或接管全部字节:
public NetworkSendResult Send(IReadOnlyList<ArraySegment<byte>> packets)
{
if (packets == null)
throw new ArgumentNullException(nameof(packets));
int totalCount = ValidateAndMeasure(packets);
bool startWorker = false;
lock (_lifecycleLock)
{
if (_disposed || _state != NetworkProviderState.Connected)
return NetworkSendResult.NotConnected;
if (totalCount > _context.SendQueueSize - _queuedSendBytes)
return NetworkSendResult.Backpressure;
if (totalCount > 0)
{
byte[] payload = CopyPackets(packets, totalCount);
_sendQueue.Enqueue(payload);
_queuedSendBytes += totalCount;
startWorker = !_sendWorkerRunning;
if (startWorker)
_sendWorkerRunning = true;
}
}
if (startWorker)
StartSendWorker();
return NetworkSendResult.Accepted;
}
private static int ValidateAndMeasure(IReadOnlyList<ArraySegment<byte>> packets)
{
long total = 0;
foreach (ArraySegment<byte> packet in packets)
{
if (packet.Array == null)
throw new ArgumentException("A segment has no backing array.", nameof(packets));
total += packet.Count;
if (total > int.MaxValue)
throw new ArgumentException("The send batch is too large.", nameof(packets));
}
return (int)total;
}
private static byte[] CopyPackets(
IReadOnlyList<ArraySegment<byte>> packets, int totalCount)
{
byte[] payload = new byte[totalCount];
int offset = 0;
foreach (ArraySegment<byte> packet in packets)
{
Buffer.BlockCopy(packet.Array, packet.Offset, payload, offset, packet.Count);
offset += packet.Count;
}
return payload;
}返回 Backpressure 时必须整批拒绝,不能发送半个 Bundle。totalCount > limit - current 还能避免先做加法产生整数溢出。
发送必须串行,避免底层 WebSocket 并发发送导致乱序:
private void StartSendWorker()
{
#if UNITY_WEBGL && !UNITY_EDITOR
_ = SendAsync();
#else
_ = System.Threading.Tasks.Task.Run(SendAsync);
#endif
}
private async System.Threading.Tasks.Task SendAsync()
{
try
{
while (true)
{
byte[] payload;
WebSocket socket;
lock (_lifecycleLock)
{
if (_disposed ||
_state != NetworkProviderState.Connected ||
_sendQueue.Count == 0)
{
_sendWorkerRunning = false;
return;
}
payload = _sendQueue.Peek();
socket = _socket;
}
if (socket == null)
throw new InvalidOperationException("The WebSocket is unavailable.");
await socket.Send(payload);
lock (_lifecycleLock)
{
if (_sendQueue.Count > 0 &&
ReferenceEquals(_sendQueue.Peek(), payload))
{
_sendQueue.Dequeue();
_queuedSendBytes -= payload.Length;
}
}
}
}
catch (Exception exception)
{
Fail("WebSocket send failed.", exception);
}
}WebGL 的浏览器 WebSocket 需要主线程亲和性;桌面 NativeWebSocket 发送路径可能阻塞,因此放到 Worker。更换底层库时必须重新确认其线程规则。
第九步:统一失败出口
连接、发送、派发和远端关闭都进入同一个失败方法:
private void Fail(string message, Exception exception)
{
WebSocket socket;
lock (_lifecycleLock)
{
if (_disposed ||
_state == NetworkProviderState.Closing ||
_state == NetworkProviderState.Closed ||
_state == NetworkProviderState.Failed)
{
return;
}
_state = NetworkProviderState.Failed;
socket = _socket;
_socket = null;
ClearSendQueueLocked();
EnqueueEvent(ProviderEvent.Closed(new NetworkCloseInfo(message, exception)));
}
Detach(socket);
BeginSocketShutdown(socket);
}先进入 Failed 再释放资源,可以阻止同时到达的 OnError、OnClose 和发送异常重复报告。关闭通知仍然排队,不能从失败发生的线程直接进入 SDK。
第十步:主动关闭和释放
主动关闭表示 SDK 已经在销毁 Session,不应再次产生 OnClosed():
public void Close() => CloseCore(false);
public void Dispose() => CloseCore(true);
private void CloseCore(bool dispose)
{
WebSocket socket;
lock (_lifecycleLock)
{
if (dispose)
_disposed = true;
if (_state == NetworkProviderState.Closed ||
_state == NetworkProviderState.Closing)
return;
if (_state == NetworkProviderState.Failed)
{
ClearEvents();
return;
}
_state = NetworkProviderState.Closing;
socket = _socket;
_socket = null;
ClearSendQueueLocked();
ClearEvents();
}
Detach(socket);
BeginSocketShutdown(socket);
_state = NetworkProviderState.Closed;
}
private void BeginSocketShutdown(WebSocket socket)
{
if (socket == null)
return;
try { socket.CancelConnection(); } catch { }
_ = CloseSocketAsync(socket);
}
private static async System.Threading.Tasks.Task CloseSocketAsync(WebSocket socket)
{
try { await socket.Close(); } catch { }
}
private void Detach(WebSocket socket)
{
if (socket == null)
return;
socket.OnOpen -= OnSocketOpen;
socket.OnMessage -= OnSocketMessage;
socket.OnError -= OnSocketError;
socket.OnClose -= OnSocketClose;
}
private void ClearEvents()
{
lock (_eventLock)
_events.Clear();
}
private void ClearSendQueueLocked()
{
_sendQueue.Clear();
_queuedSendBytes = 0;
_sendWorkerRunning = false;
}Close() 和 Dispose() 必须幂等,因为 SDK 的确定性清理会依次尝试二者。事件解绑还能避免第三方 Socket 长期持有 Provider,或旧连接在新 Session 建立后继续回调。
第十一步:注册 Provider
using System.Collections.Generic;
using KBEngine;
using KBEngine.UnityWebSocket;
public sealed class ClientApp : UnityKBEMain
{
private readonly Dictionary<string, string> _domainMapping =
new Dictionary<string, string>
{
{ "192.168.36.128", "wss.kbelab.com" },
};
private readonly Dictionary<int, int> _portMapping =
new Dictionary<int, int>
{
{ 20013, 443 },
{ 20015, 444 },
};
protected override INetworkProviderFactory CreateCustomNetworkProviderFactory()
{
if (networkType != KBEngineApp.NETWORK_TYPE.CUSTOM &&
networkType != KBEngineApp.NETWORK_TYPE.CUSTOM_ALL)
return null;
return new UnityWebSocketProviderFactory(
secure: true,
path: "/",
domainMapping: _domainMapping,
portMapping: _portMapping);
}
}WebGL、微信小游戏、抖音小游戏等不能使用普通 Loginapp TCP 的运行时应选择 CUSTOM_ALL。只有 Loginapp 可以继续使用 SDK TCP、而 Baseapp 使用自定义传输时,才选择 CUSTOM。
替换为其他网络 API
适配小游戏或引擎自带 Socket 时,保留工厂、状态机、事件队列、发送所有权、背压和幂等关闭,只替换下表中的传输调用:
| WebSocket 示例 | 其他平台对应能力 |
|---|---|
new WebSocket(url) | 创建平台或引擎 Socket |
socket.Connect() | 发起异步连接 |
OnOpen | 连接成功回调 |
OnMessage(byte[]) | 二进制接收回调 |
socket.Send(byte[]) | 发送二进制数据 |
OnError / OnClose | 错误和远端关闭回调 |
CancelConnection() / Close() | 取消连接并释放句柄 |
微信小游戏、抖音小游戏回调可能提供 ArrayBuffer 或平台自有 Buffer。转换后的内存必须在 Process() 派发期间仍有效。客户端引擎自带 Socket 若已经管理线程、日志和缓冲池,Provider 应复用这些能力,不要再实现一套重复的网络线程。
性能与并发
CPU、内存和 GC
示例每个发送批次执行一次 O(n) 合并复制,并分配一个 byte[]。高频小包可使用 ArrayPool<byte>,但只能在异步发送完成或队列清理后归还,并需要保存有效长度。过早归还会造成随机协议损坏。WebSocket 必须发送二进制数据,避免 Base64 和文本编解码。
消息数量和网络开销
不要跨多个 SDK Send() 无限合批,这会提高延迟和内存。若确实需要合批,应同时设置字节数、消息数和等待时间上限,并严格保持字节顺序。
锁竞争
锁内只做状态检查、引用替换和队列操作。实际连接、发送、关闭和 Listener 回调放在锁外,避免底层库同步回调时重入死锁。
Tick 压力
网络突发可能使 DrainEvents() 单 Tick 执行过久。生产实现可加入事件数量或耗时预算,并监控事件队列长度、发送排队字节、单 Tick 派发数和最长排队时间。预算只能延后处理,不能丢弃或重排事件。
测试清单
功能
CUSTOM_ALL完成 Loginapp、Baseapp 切换和 Entity 创建。CUSTOM完成 TCP Loginapp 和 WebSocket Baseapp 混合链路。- Relogin 创建新 Provider,并使用保存的 Baseapp Endpoint。
- 半条消息、多条合并消息和跨多次接收的数据都能解析。
- WSS、路径、查询参数、域名和端口映射正确。
生命周期与并发
- 连接期间主动退出。
- 已连接时远端关闭。
OnError与OnClose连续到达,只报告一次失败。- 发送过程中关闭,不重复释放或回调。
- 多次调用
Close()、Dispose()不抛异常。
压力和稳定性
- 发送速度高于出口时,在队列上限返回
Backpressure。 - 大包接近
SEND_QUEUE_MAX时不溢出、不发送半包。 - 长时间运行时排队字节、事件队列和托管内存不持续增长。
- 主线程卡顿恢复后事件仍保持顺序且不重复。
- Standalone、WebGL 和目标小游戏真机分别验证。
常见错误
- 平台回调直接调用 Listener:会让 SDK 协议状态运行在不可控线程,应改为排队并从
Process()派发。 - 保存 SDK 传入的 ArraySegment:底层数组可能立即被对象池复用,必须先复制或接管。
- 工厂只创建一个 Provider:旧实例已处于终态,不能用于 Baseapp 或 Relogin。
- 把 WebSocket Message 当成 KBEngine 消息:传输边界不等于协议边界,Provider 只保证字节顺序。
- 主动关闭仍报告 OnClosed:会造成重复断线事件,只有远端关闭和不可恢复错误才通知。
- 无限增大发送队列:这只会把网络背压变成内存增长和高延迟,应从消息频率和出口带宽解决。
完成标准
一个 Provider 只有同时满足以下条件才算完成:
- 向 SDK 提供有序、未经修改的字节流。
- 所有 Listener 回调只发生在
Process()。 - 返回
Accepted后不再引用 SDK 对象池内存。 - 连接、发送、远端关闭和主动关闭具有确定状态。
- 发送和事件处理有明确资源上限或 Tick 预算。
- Loginapp、Baseapp 和 Relogin 使用互相隔离的新实例。
- 关闭与释放幂等,不泄漏线程、句柄、回调或缓冲。
- 目标平台真机验证通过,而不只是 Editor 中可用。
完成实现后,回到《网络层 Provider》继续检查 WSS、反向代理、线程模式和平台部署。
