OpenRobot.Framework.Communication 1.2.0

There is a newer version of this package available.
See the version list below for details.
dotnet add package OpenRobot.Framework.Communication --version 1.2.0
                    
NuGet\Install-Package OpenRobot.Framework.Communication -Version 1.2.0
                    
This command is intended to be used within the Package Manager Console in Visual Studio, as it uses the NuGet module's version of Install-Package.
<PackageReference Include="OpenRobot.Framework.Communication" Version="1.2.0" />
                    
For projects that support PackageReference, copy this XML node into the project file to reference the package.
<PackageVersion Include="OpenRobot.Framework.Communication" Version="1.2.0" />
                    
Directory.Packages.props
<PackageReference Include="OpenRobot.Framework.Communication" />
                    
Project file
For projects that support Central Package Management (CPM), copy this XML node into the solution Directory.Packages.props file to version the package.
paket add OpenRobot.Framework.Communication --version 1.2.0
                    
#r "nuget: OpenRobot.Framework.Communication, 1.2.0"
                    
#r directive can be used in F# Interactive and Polyglot Notebooks. Copy this into the interactive tool or source code of the script to reference the package.
#:package OpenRobot.Framework.Communication@1.2.0
                    
#:package directive can be used in C# file-based apps starting in .NET 10 preview 4. Copy this into a .cs file before any lines of code to reference the package.
#addin nuget:?package=OpenRobot.Framework.Communication&version=1.2.0
                    
Install as a Cake Addin
#tool nuget:?package=OpenRobot.Framework.Communication&version=1.2.0
                    
Install as a Cake Tool

OpenRobot.Framework.Communication

MQTT 连接、IPC 跨进程通讯与 SignalR 长连接通用库(net8.0,单包内含 Mqtt / Ipc / Signalr 三个独立命名空间)。

  • MQTT:连接控制 + 自动重连 + 重连后重订阅 + 按序补发;订阅以 (Topic, Handler) 批量注册,收包经 Channel 消费队列派发;QoS / Retain / MQTT5 ResponseTopic 与 CorrelationData 可配。
  • IPC:一个服务端 + 多个客户端。服务端按握手 ClientId 维护会话表,定向发送 / 广播 / 单会话断开互不影响;Client 与 Session 接口同构,双向「单向发送(不等待返回)」与「请求-响应(带超时)」都可用;传输可在命名管道 / TCP 之间切换,也可用 MultiTransport 让同一个服务端同时监听多种传输(两种传输上的客户端在同一张会话表里,可互相路由)。
  • SignalR:连接控制 + 自动重连 + 处理器跨重连存活;普通调用(异常 / Result 两种模式)、服务端→客户端 流式接收(Pull + Push 广播)、客户端→服务端 流式上传(含 byte[] 自动分片);空闲超时看门狗(长时间收不到业务数据可配置重连 / 停止 / 失败);授权令牌与自定义头预留。
  • 协议栈类型(MQTTnet、Microsoft.AspNetCore.SignalR.Client)不出现在公共 API:换掉协议栈不影响调用方。
  • 日志全部经 ICommunicationLogger,宿主可实现后接管输出(也支持直接用 Microsoft.Extensions.Logging)。
  • 无副作用:不写环境变量、不落盘、不生成持久状态;凭据只读,全部由宿主提供。

安装

<PackageReference Include="OpenRobot.Framework.Communication" Version="1.2.0" />

包由本仓库产出:dotnet pack common/OpenRobot.Framework.Communication/OpenRobot.Framework.Communication.csproj -c Release, 产物为 OpenRobot.Framework.Communication.1.2.0.nupkg(lib/net8.0 单个程序集)与配套 .snupkg。

当前版本 1.2.0(新增 IPC 信封协议层)。版本号在 OpenRobot.Framework.Communication.csproj 的 <Version> 中维护, 后续发布会在此处一并升级(含上文引用与 版本历史 一节)。

依赖:

  • MQTTnet 4.3.7.1207(仅内部使用)
  • Microsoft.AspNetCore.SignalR.Client 8.0.8 + Microsoft.AspNetCore.SignalR.Protocols.MessagePack 8.0.8(仅内部使用)
  • Microsoft.Extensions.Logging.Abstractions 8.0.2
  • Microsoft.Extensions.DependencyInjection.Abstractions 8.0.2

快速开始

MQTT

using OpenRobot.Framework.Communication;
using OpenRobot.Framework.Communication.Mqtt;

var mqtt = new MqttClient(new MqttOptions
{
    Name = "primary",
    Broker = "127.0.0.1",
    Port = 1883,
    Logger = new ConsoleCommunicationLogger(),   // 或实现 ICommunicationLogger
});

mqtt.AddSubscriptions(
    // 处理完要回执:返回值非 null 即回执(带 CancellationToken 的重载用 (m, ct) => ...)
    new MqttSubscription("ota/down/#", HandleOtaDownAsync),
    new MqttSubscription("config/changed", msg => Console.WriteLine(msg.Payload)));

await mqtt.StartAsync();     // 非阻塞:见下方「StartAsync 不等首连」
await Task.Delay(500);
await mqtt.PublishAsync("ota/up", new { deviceId = "d1", state = "ok" });

static Task<string?> HandleOtaDownAsync(MqttMessage msg, CancellationToken ct)
    => Task.FromResult<string?>("{\"accepted\":true}");

IPC 服务端(一个端点,多个客户端)

using OpenRobot.Framework.Communication.Ipc;

await using var service = new IpcService(new IpcOptions { Name = "hub", PipePrefix = "EdgeControl" });

service.On("telemetry.report", inbound =>
{
    // inbound.SourceId 是发来这条消息的客户端 ClientId
    Console.WriteLine($"{inbound.SourceId}: {inbound.Message.Payload}");
    return Task.CompletedTask;
});
// 需要强类型载荷时显式指定泛型参数(否则 (payload, ctx) 双参 lambda 无法推断 T)
service.On<TelemetryReport>("telemetry.report", report => HandleReportAsync(report));

service.OnRequest("api.proxy.request", async ctx =>
{
    if (!await IsAuthorizedAsync(ctx.SourceId)) { await ctx.ReplyErrorAsync("unauthorized"); return; }
    await ctx.ReplyAsync(new { ok = true });
});

service.ClientConnected += (_, e) => Console.WriteLine($"上线: {e.Session.ClientId} @ {e.Session.RemoteEndpoint}");

await service.StartAsync();                       // 返回时已可连接
await service.SendToAsync("upgrader", "ota.dispatch", new { url = "https://..." });
await service.BroadcastAsync("config.changed", new { revision = 42 });
var reply = await service.RequestAsync("upgrader", "api.proxy.request", new { path = "/x" },
    TimeSpan.FromSeconds(30));

IPC 客户端(独立进程)

await using var client = new IpcClient(new IpcOptions
{
    Name = "upgrader",
    ClientId = "upgrader",        // 服务端按它区分会话
    PipePrefix = "EdgeControl",
    DefaultRequestTimeout = TimeSpan.FromSeconds(60),
});

client.On<OtaDispatch>("ota.dispatch", async task => await RunAsync(task));   // 不关心来源
await client.StartAsync();

// 客户端也能反向往服务端发请求(双向对称)
var cfg = await client.RequestAsync("config.get", new { key = "revision" });

IPC 服务端同时监听多种传输(MultiTransport)

一个 IpcService 默认只监听一种传输。要让同一个服务端同时接受命名管道与 TCP 的客户端, 并让它们互相路由(走 TCP 的 A 把消息发给走管道的 C),用 MultiTransport 把多个监听端点合成一个:

using OpenRobot.Framework.Communication.Ipc;
using OpenRobot.Framework.Communication.Ipc.Transport;

await using var service = new IpcService(new IpcOptions
{
    Name = "hub",
    TransportOverride = new MultiTransport(
        new NamedPipeTransport("EdgeControl.hub", maxInstances: 128),
        new TcpTransport("0.0.0.0", 9000)),
});

service.On<Telemetry>("telemetry.report", async (inbound, payload) =>
{
    // 会话表只有一张:只按 ClientId 路由,不关心对端走的是哪种传输
    await service.SendToAsync(payload.TargetClientId, "telemetry.forward", payload);
});

await service.StartAsync();   // 返回时两个端点都已可连接
// 之后 SendToAsync / BroadcastAsync / RequestAsync 对两种传输上的客户端一视同仁

客户端不用改 —— 各自按 Transport 选自己要用的那一种即可。MultiTransport 是服务端专用, 在客户端用它(ConnectAsync)会抛 NotSupportedException:连哪一个属于宿主业务决策, 组合器若自行挑选会让重连在多个端点间来回跳、故障现场无法复现。

SignalR

using OpenRobot.Framework.Communication.Signalr;

await using var hub = new SignalrClient(new SignalrOptions
{
    Name = "shell",
    HubUrl = "http://127.0.0.1:9900/signalr-remote-shell",
    LogTag = "Signalr.shell",
    Protocol = SignalrProtocolKind.MessagePack,   // 默认;两端必须一致
    ReconnectDelaysMs = [0, 2_000, 10_000, 30_000],   // 序列用尽后复用最后一项
    MaxReconnectAttempts = 0,                          // 0 = 真的无限
    IdleTimeoutMs = 60_000,                            // 0 = 不启用(默认)
    IdleAction = SignalrIdleAction.Reconnect,
});

// ① 服务端 → 客户端 事件:处理器跨重连存活
hub.On<string>("ShellOutput", text => { Console.WriteLine(text); return Task.CompletedTask; });

await hub.StartAsync();     // 非阻塞:只启动后台重连循环,不等首连结果

// ② 普通调用(异常模式 / Result 模式)
await hub.InvokeAsync("ResizeTerminal", sessionId, 120, 40);
var result = await hub.TryInvokeAsync<string>("GetVersion");   // 不抛,返回 Success/ErrorType

// ③ 服务端 → 客户端 流(Pull + Push 广播并存)
var stream = await hub.StreamAsync<string>("StreamShellOutput", ct, sessionId);
_ = Task.Run(async () =>
{
    await foreach (var chunk in stream.ReadAllAsync()) Console.Write(chunk);   // Pull
});
using var push = stream.Subscribe(chunk => Send(chunk));                       // Push,各自收全量

// ④ 客户端 → 服务端 流:任意长度 byte[] 自动按 16 KiB 分片
var upload = hub.CreateBinaryUploadStream("ReceiveShellOutput");
await upload.WriteAsync(rawBytes);   // 任意长度,内部切片
upload.Complete();                   // 或 upload.WriteAsync 后自然完成
await upload.WaitForCompletionAsync();

IPC 信封协议(OpenRobot.Framework.Communication.Ipc.Envelope)

IPC 帧只带一个固定类型(IpcWire.FrameType = "edge.message"),业务的一切都在信封里。

// 出方向:把业务载荷装进信封发出去
public sealed class MyHandler : IIpcMessageHandler
{
    private readonly IIpcEnvelopeSender _sender;   // ⚠️ 注 sender,不要注 gateway(会成 DI 环)

    public string Topic => "my.topic";

    public async Task HandleAsync(IpcMessageContext context, CancellationToken ct)
    {
        if (context.IsRequest) await context.ReplyAsync(new { ok = true }, ct);
        await _sender.SendAsync("other.topic", new { hello = "world" }, ct: ct);
    }
}
  • IpcEnvelope —— 五字段(topic / target / source / error / payload), 字段名 camelCase 是契约,服务端另有一份同源副本,改名会静默断开所有客户端。 业务 ID、业务码、业务错误说明全在 payload 里;error 只承载到不了处理器的失败
  • IpcWire —— 线格式常量(FrameType)
  • IIpcMessageHandler —— 入方向扩展点,一个 topic 一个实现类
  • IIpcConnectionHook —— 出方向连接钩子,每次连上后触发,须幂等
  • IIpcEnvelopeSender / IpcEnvelopeSender —— 出方向门面(单向 SendAsync + 请求-响应 RequestAsync / RequestToAsync,超时与未连接一律返回 null 不抛)
  • IIpcEnvelopeGateway / IpcEnvelopeGateway —— 入方向网关。一个进程只能注册一个: 库的 On 按帧类型建表、后注册者覆盖前者,网关独占注册后再按信封的 topic 二次分发。 它不继承 IIpcEnvelopeSender —— 想发消息就注入 sender
  • IpcMessageContext —— 处理器上下文。回执关联只认帧的 RequestId,与信封无关; 单向消息上调 ReplyAsync 是 no-op
  • 本机名(信封的 source)由 Attach 从 client.ClientId 取,不从配置取 —— 与握手报给服务端的是同一个值。构造函数因此不需要任何配置
  • ⚠️ 信封层固定用 IpcJson.Default 序列化载荷,不受 IpcOptions.Json 影响 —— IpcOptions.Json 只管 IPC 帧的序列化。别以为配了自定义转换器 (例如给 DateTime 换格式)就能让信封的 payload 跟着变:帧与信封会各按各的来

依赖方向:网关 → 处理器 → IIpcEnvelopeSender。业务处理器注入的必须是 sender —— 注 gateway 会让 DI 环在类型上成立,启动时才崩。(1.2.0 起 IIpcEnvelopeGateway 不再继承发送器,所以那种写法在 1.2.0+ 是编译错误,直接写不出来。)


API 速查

IMqttClient

成员 说明
Name 连接名(MqttOptions.Name,多连接区分日志用)
IsConnected / State 是否已连接 / MqttConnectionState(Disconnected、Connecting、Connected、Reconnecting、Failed、Stopped)
ConnectionStateChanged 状态转换事件,一次转换触发一次
StartAsync(ct) 非阻塞启动 + 后台自动重连循环,幂等
ConnectAsync(ct) 单次连接(不等自动重连);失败会抛
StopAsync() 停止自动重连并断开,保留订阅注册,可再次 StartAsync
AddSubscriptions(...) 批量注册订阅,按 TopicFilter 去重(同 filter 只保留首个)
SubscribeAsync(filter, ct) 只确保订阅、不挂处理逻辑
UnsubscribeAsync(filter, ct) 从注册表摘除 + 移出重订阅清单(不发 UNSUBSCRIBE,见「已知限制」)
AddInterceptor(i) 注册收包前置拦截器(入队前按序同步执行)
PublishAsync(topic, payload[, options]) 未连接或发布失败时入补发队列,重连后按序补发

MqttSubscription 四个构造重载:Func<MqttMessage, CancellationToken, Task<string?>> 与 Func<MqttMessage, Task<string?>>(返回值非 null 且消息带 ResponseTopic 时自动回执)、Func<MqttMessage, Task>、Action<MqttMessage>(后两者不回复)。

带 Task<string?> 的两个重载必须都在:只有三参版时,一参 lambda m => Task.FromResult<string?>(json) 会因元数对不上而静默绑定到「不回执」的 Func<MqttMessage, Task>(Task<string?> 可隐式转 Task),回执凭空消失且不报错。

IIpcService(服务端)

成员 说明
Name / IsListening / ClientIds 端点名 / 是否监听中 / 在线客户端 ID 快照
ClientConnected / ClientDisconnected 会话建立 / 断开(含被新连接顶掉,每会话恰好一次)
StartAsync(ct) / StopAsync() 返回时端点已可连接;停止会断开全部会话(保留处理器,可重启)
GetSession(id) / GetSessions() 按 ClientId 取会话 / 取全部会话快照
SendToAsync(id, ...) 定向单向发送;客户端不在线返回 false(不缓存)
BroadcastAsync(...) 广播给全部在线会话;离线跳过,单会话失败不影响其余
RequestAsync(id, ...) 定向请求-响应;不在线或超时返回 null(不抛)
On / OnRequest(继承 IIpcMessageHandlers) 服务级处理器,对所有会话生效(含之后才连上的客户端);会话级同名处理器优先

IIpcClient / IIpcSession(共用 IIpcEndpoint)

成员 说明
Name Client 侧 = 自己的 ClientId;Session 侧 = 该客户端的 ClientId
IsConnected / State / ConnectionStateChanged 连接状态与状态事件(客户端每次重连成功都会触发)
SendAsync(message) / SendAsync<T>(type, payload) 单向发送:仅入队即返回(ValueTask),队列满丢最旧 + Warn;未连接时保留、连上后按序发出
RequestAsync(...) 请求-响应;超时或失败返回 null(不抛)
On / OnRequest 注册处理器(Session 侧仅对该客户端生效,且优先于服务级)
ClientId / ReconnectAttempts 仅 IIpcClient:握手标识 / 当前连续重连次数
RemoteEndpoint / ConnectedAt / Hello 仅 IIpcSession:对端地址 / 握手完成时间 / 握手载荷

处理器内可用 IpcInboundMessage.SourceId 判断来源;请求处理器用 IpcRequestContext.ReplyAsync / ReplyErrorAsync 回传(ReplyErrorAsync 回 {"success":false,"error":...},请求方快速失败而非等超时)。

__ipc. 前缀由库独占(握手帧 __ipc.hello),业务消息类型不得占用。

ISignalrClient

成员 说明
Name 连接名(SignalrOptions.Name,多连接区分日志用)
IsConnected / State / ConnectionId 是否已连接 / SignalrConnectionState(Disconnected、Connecting、Connected、Reconnecting、Failed、Stopped)/ 底层连接 ID
ConnectionStateChanged 状态转换事件,一次真实转换触发一次
TimeSinceLastActivity 距最近一次收到业务数据(On<T> 回调 / 流数据项 / Invoke 响应)的时长
StartAsync(ct) 非阻塞启动 + 后台自动重连循环,幂等
ConnectAsync(ct) 单次连接(连上才返回、失败抛出);不启动自动重连
StopAsync() 停循环并断开,保留处理器与配置,可再次 StartAsync
DisconnectAsync() 只断开、保留处理器(循环会在退避后重新连上)
InvokeAsync(...) ×4 普通调用,失败抛 SignalrConnectionException;用户主动取消原样抛 OperationCanceledException
TryInvokeAsync(...) ×2 Result 模式,业务失败不抛,返回 SignalrInvokeResult(含 ErrorType)
On<T1..T8>(method, handler) 注册服务端→客户端 处理器,返回的 IDisposable 只摘除这一个
RemoveHandler(method) / RemoveAllHandlers() 按方法名移除全部 / 清空全部
StreamAsync<T>(method, ct, args) 发起服务端→客户端 流,返回可 Pull / Push 的订阅(断线自动重建;懒建立,见下)
CreateUploadStream<T>(method, ct) 创建客户端→服务端 上传流(每个写入项 = 线上一个元素)
CreateBinaryUploadStream(method, chunkSize, ct) 客户端流式 byte[]:写入任意长度,内部按分片切开
IdleTimedOut 空闲超时事件(已按 IdleAction 处理完之后触发)

错误分类 SignalrErrorType:NotConnected / Timeout / ConnectionLost / ServerError / Unknown。

重连权威只有一处:本库不配置 SignalR 内建的 WithAutomaticReconnect,退避与尝试次数全部由 SignalrClient 的后台循环独占决策。因此 MaxReconnectAttempts 是真实计数,ReconnectDelaysMs 序列用尽后复用最后一项(不会像原生数组重载那样静默停止重连)。

StartAsync 不等首连:与 IMqttClient 同形。需要「连上再继续」请用 ConnectAsync,或订阅 ConnectionStateChanged 等 Connected。

StreamAsync 会等:连接处于 Connecting / Reconnecting / Disconnected 时等待连接就绪(「还没连上」与「连不上」不是一回事),只有进入 Failed / Stopped 才立即抛 NotConnected。ct 是等待上限的唯一手段。

流是懒建立的:StreamAsync 只登记订阅,首个消费者接入(Subscribe / ReadAllAsync)后才真的向服务端建流 —— 这样「建好流到消费者挂上」之间到达的数据不会静默丢失。代价是后接入的消费者只收得到从它接入那一刻起的数据;也意味着没有任何消费者时该订阅既不建立流、也不会完成(连接进入终态才结束)。

DI 集成

services.AddOpenRobotCommunication(o =>
{
    o.Logger = myLogger;                                   // 全局日志(可选)
    o.AddMqtt("primary", m => { m.Broker = "broker.local"; m.LogTag = "Mqtt.primary"; });
    o.AddIpcService("hub", i => { i.Name = "hub"; i.PipePrefix = "EdgeControl"; i.MaxClients = 32; });
    o.AddIpcClient("upgrader", i => { i.Name = "upgrader"; i.ClientId = "upgrader"; });
    o.AddSignalr("shell", s => { s.HubUrl = "http://127.0.0.1:9900/signalr-remote-shell"; s.LogTag = "Signalr.shell"; });
});

var client = sp.GetRequiredKeyedService<IMqttClient>("primary");
await client.StartAsync();
  • 每个 MQTT 连接 = keyed singleton IMqttClient;每个 IPC 端点 = keyed singleton IIpcClient / IIpcService;每个 SignalR 连接 = keyed singleton ISignalrClient。
  • 不自动连接 / 监听:注册只构造实例,由宿主在适当时机调用 StartAsync。
  • 同名重复登记后者覆盖前者:三种连接都按名索引,后登记的配置会替换先前的(不报错、不追加)。
  • 日志优先级:MqttOptions/IpcOptions/SignalrOptions.Logger > CommunicationOptions.Logger > 容器的 ILoggerFactory > 静默。
  • 容器必须异步释放:四个服务都是 IAsyncDisposable 且不实现 IDisposable,同步 ServiceProvider.Dispose() 会抛 InvalidOperationException("type only implements IAsyncDisposable. Use DisposeAsync to dispose the container.")。用 await using var sp = services.BuildServiceProvider();。

配置项

MqttOptions

项 默认 说明
Name "default" 连接名(日志 + MqttMessage.ClientName)
LogTag null → "Mqtt" 日志 tag
Broker / Port "" / 1883 broker 地址与端口
ClientId / Username / Password null 显式凭据(非空即生效,优先于 CredentialProvider)
KeepAliveSeconds 60 心跳间隔
ConnectTimeoutMs 10000 连接建立超时
UseTls / AllowUntrustedCertificate false / false TLS 开关与跳过证书校验(内网用,生产应关闭)
DefaultQuality AtMostOnce 发布默认 QoS
CleanSession true 不依赖 broker 持久会话,重连靠客户端侧重订阅 + 补发
ReconnectDelayMs 5000 重连检测间隔
MaxReconnectAttempts 0 0 = 不限;达上限置 Failed 并停止重连
PendingPublishCapacity 10000 补发队列上限(超限丢最旧 + Warn)
PendingPublishMaxRetry 3 单条最大补发次数(超限丢弃 + Error)
DispatcherQueueCapacity 10000 消费队列上限(满丢最旧 + Warn)
DispatcherConsumerCount 1 1 保证同连接内按到达顺序处理;>1 提速但放弃顺序保证
CredentialProvider null IMqttCredentialProvider,显式凭据未配置时调用
Logger null 日志实现(null 静默)

IpcOptions(Service 与 Client 共用)

项 默认 说明
Name "" Service 为监听名;Client 为默认标识名(见 ClientId)
LogTag null → "Ipc" 日志 tag
Transport NamedPipe NamedPipe(Windows 管道 / Linux UDS)或 Tcp
PipePrefix "OpenRobot" 管道名 / socket 文件前缀
TcpHost / TcpPort "127.0.0.1" / 0 Transport = Tcp 时生效
ClientId null Client 角色专用:为空取 Name,仍为空则进程内生成 Guid
ReconnectDelayMs 3000 Client 重连间隔
ConnectTimeoutMs 10000 连接建立超时
HandshakeTimeoutMs 10000 Service 等待 hello 首帧的时限,超时关连接
MaxClients 128 Service 最大并发会话数(超出拒绝新连接 + Warn)。命名管道下实际封顶 252
OutboundQueueCapacity 10000 单向发送队列上限(满丢最旧 + Warn)
DefaultRequestTimeout 30s RequestAsync 未显式指定超时时的默认值
MaxFrameBytes 64 MB 单帧最大字节数(防御异常长度)
Json null → IpcJson.Default 载荷序列化配置(冻结为只读,运行期不可改)
TransportOverride null 自定义 IIpcTransport(测试 / 自定义传输 / 服务端多传输 MultiTransport;非 null 时忽略上面四个传输项)
Logger null 日志实现(null 静默)

写失败一律拆连接、不在原连接上重试:写失败可能已在线上留下半截帧,重试会把新帧追加进半截帧导致对端静默解析错乱。在途消息保留,由下一条连接优先补发。

SignalrOptions

项 默认 说明
Name "default" 连接名(多连接日志区分)
LogTag null → "Signalr" 日志 tag
HubUrl "" Hub 完整地址(含路径,如 http://host:9900/signalr-remote-shell);须为 http/https 绝对地址
ConnectTimeoutMs 30000 连接建立超时
HandshakeTimeout 15s SignalR 握手阶段超时
KeepAliveInterval 15s 心跳间隔
ServerTimeout 30s 传输层服务端超时;构建连接时校验 须 ≥ 2 × KeepAliveInterval
Protocol MessagePack MessagePack / Json;两端必须一致,否则握手被拒
ConfigureProtocolServices null 协议高级配置逃生舱(在 AddJsonProtocol/AddMessagePackProtocol 之后执行,Action<IServiceCollection> 形式,公共 API 不泄漏协议类型)
AccessTokenProvider null 授权预留:每次协商 / 重连时调用取最新令牌;库只读取,不获取 / 刷新 / 持久化
Headers null 授权预留:随请求发送的自定义头(如 Abp.TenantId)
AutoReconnect true 断线后是否自动重连
ReconnectDelaysMs [0, 2000, 10000, 30000] 第 N 次重连前等第 N 项;序列用尽后复用最后一项
MaxReconnectAttempts 0 0 = 无限;>0 时达上限置 Failed 并停止重连
IdleTimeoutMs 0 0 = 不启用(opt-in):连接期间超过该时长未收到业务数据则触发 IdleAction
IdleAction Reconnect Reconnect / Stop / Fail
IdleCheckIntervalMs 1000 看门狗检查周期(仅影响检测精度)
IdleProbe null 可选探活钩子:超时后先调用,返回 true 视为健康(重置计时、不动作)
StreamBufferCapacity 10000 服务端→客户端 流的每订阅者缓冲上限(满丢最旧 + Warn)
UploadBufferCapacity 10000 客户端→服务端 上传缓冲上限
UploadBufferMode DropOldest 上传缓冲满时策略(DropOldest / DropNewest / Wait)
BinaryChunkSize 16 KiB 二进制上传默认分片大小(见下方说明)
StreamRetryDelayMs 1000 流中断后重新建流前的等待
Logger null 日志实现(null 静默)

空闲超时的判定口径:「业务数据」= On<T> 事件回调、流式订阅数据项、Invoke 响应返回、连接建立;不含传输层心跳——心跳正常但业务静默正是它要抓的状态。副作用是:宿主没为某类服务端推送注册 On<T> 时,SignalR 客户端会静默丢弃那些消息,本库无从计为活动(「宿主没订阅的数据 = 没收到」)。不该断开就别开本项。

为什么 BinaryChunkSize 默认 16 KiB 而不是 64 KiB:SignalR 服务端的 HubOptions.MaximumReceiveMessageSize 默认只有 32 KB,超过即拒收该消息并断开连接。上行流式分片的每一片是独立的 Hub 消息,64 KiB 分片会被默认配置的服务端直接拒收。服务端若已放宽该限制可自行调大。本项只影响上行分片——MaximumReceiveMessageSize 只存在于服务端配置中,客户端程序集里没有这个设置。

授权为预留能力:本库不实现令牌获取 / 刷新 / 持久化,只提供 AccessTokenProvider 与 Headers 两个注入口,由宿主提供最新值。两条已知限制:WebSocket 长连接建立后不会刷新令牌(只在下一次协商 / 重连时刷新);SignalR 只在协商请求上带自定义头,WebSocket 握手走查询串(令牌除外)——需要全程生效的信息请放进令牌或 Hub 方法参数。


平台说明

  • Windows 命名管道:端点名为 {PipePrefix}.{Name}(如 EdgeControl.hub),保留原大小写;管道名不区分大小写。同一服务端的管道实例数受系统上限 254 约束(库会为「在途 Accept」与「待拒绝连接」各留 1 个实例,故会话数实际封顶 252)。需要更多并发会话请改用 Transport = Tcp。
  • Linux / macOS:System.IO.Pipes 自动映射为 Unix Domain Socket,端点名为 /tmp/{prefix}.{name}.sock(转小写——UDS 区分大小写,且路径总长受系统上限约 108 字节约束)。PipePrefix 与 Name 应保持简短,过长会导致连接失败。
  • TCP:跨机器可用,但本期为明文传输,无 TLS。仅建议在受信内网使用;需要跨不可信网络时应在链路上自行加隧道 / 加密,或等待后续 TLS 扩展点。
  • 多传输并存:MultiTransport(服务端专用)把多个监听端点合成一个 IIpcTransport,经 TransportOverride 注入即可 —— 它不改动 IpcService,因为会话层只认 Stream,本来就不关心流来自哪种传输。各传输共用同一张会话表,故 SendToAsync / BroadcastAsync / RequestAsync 天然跨传输可用。
  • 传输可测性:IIpcTransport 是接缝,InMemoryTransport 用于单进程内的确定性测试(TransportOverride 注入)。注意它是 1:1 配对,多次 AcceptAsync 会返回同一组通道,建模不出多条独立连接,多客户端用例请用真实传输。

已知限制

  1. UnsubscribeAsync 不发 UNSUBSCRIBE:内部 broker 接缝(IMqttBrokerConnection)只提供 Connect / Disconnect / Publish / Subscribe,没有 Unsubscribe。因此"注销订阅"只做「从注册表摘除 + 移出重订阅清单」——broker 可能仍在投递,但注册表中已无匹配项,派发器只记 Warn 而不调用任何处理器。需要真实退订需先扩接缝。
  2. MqttClient.StartAsync 不等首次连接:它只启动后台重连循环,首连失败不向上抛。await StartAsync() 返回后 IsConnected 仍可能为 false。需要「连上再继续」请用 ConnectAsync,或订阅 ConnectionStateChanged 等 Connected。宿主的启动就绪判断必须基于 State / ConnectionStateChanged。
  3. 容器需异步释放:见「DI 集成」一节。
  4. MQTT 消息顺序仅在 DispatcherConsumerCount = 1(默认)时保证;提高并发数即放弃该保证。
  5. IPC 单向发送不保证送达:SendAsync 只入队即返回,队列满时丢最旧(记 Warn);需要确认送达请用 RequestAsync。
  6. IPC 服务端不缓存离线消息:SendToAsync 对离线客户端返回 false;BroadcastAsync 跳过离线会话。
  7. IPC 握手不可省略:客户端连接后第一帧必须是 __ipc.hello,首帧非 hello / ClientId 为空 / 帧非法 / 超 HandshakeTimeoutMs 的连接一律被关闭且不建会话。
  8. 凭据只读:库不生成、不持久化、不写环境变量。ClientId 三级回退中最后一级是进程内随机 Guid,多实例场景应显式配置 ClientId。
  9. MqttClient 公开构造只有 (MqttOptions, IMqttCredentialProvider?):协议栈工厂与 logger 的注入构造为 internal(供同程序集与测试使用),宿主无法替换协议栈实现。
  10. SignalrClient.StartAsync 不等首次连接:同 MQTT。await StartAsync() 返回后 IsConnected 仍可能为 false;State 会依次走 Connecting → Connected(首连失败则一直重试)。
  11. StreamAsync 会等待连接就绪:未 StartAsync 又无超时的客户端上建流会一直等,需要「快速失败」请传带超时的 ct。
  12. 流是懒建立的:首个消费者接入后才向服务端建流(避免接入前到达的数据无人接住);后接入的消费者只收到接入之后的数据,且没有消费者时订阅不会建立流也不会完成。
  13. 重连后服务端流从头开始:流订阅断线会自动重建,但服务端会重新执行流方法,产生重复项。库不做内建去重,宿主按业务序号自行处理。
  14. 上传流不保证 at-least-once:缓冲区内的项重连后继续发送(至少一次),已取走在途的项断线即丢失、不重放(至多一次)。真 at-least-once 需要「未确认项队列 + 服务端 ack」的跨端协议。另外重连对服务端而言是一次全新的 Hub 调用(新会话),需服务端自行做幂等。
  15. 客户端无 MaximumReceiveMessageSize:该限制只存在于服务端 HubOptions(默认 32 KB)。下行无等价硬限制;上行分片受服务端限制约束,见 BinaryChunkSize。
  16. 空闲超时把「没订阅的数据」视作「没收到」:未注册 On<T> 的服务端推送被 SignalR 客户端静默丢弃,本库无从计为活动(见 SignalrOptions 一节)。
  17. ISignalrHubConnectionFactory 是 internal 接缝:测试 / 自定义连接构建经 InternalsVisibleTo 使用,宿主无法替换 Hub 连接实现。SignalrOptions.ConnectionOverride 同理。
  18. MultiTransport 是服务端专用:ConnectAsync 抛 NotSupportedException。客户端「连哪一个」由宿主决定,组合器自行挑选会让重连在端点间乱跳。
  19. MultiTransport 下 MaxClients 是所有传输的合计上限:各传输共用同一张会话表,故它是「两个端点加起来最多 N 个会话」,不是每个端点各 N 个。
  20. MultiTransport 下无法从会话反查来源传输:IIpcSession.RemoteEndpoint 是本端监听端点的拼接串(如 EdgeControl.hub | 0.0.0.0:9000),不是对端地址 —— 这一点在单传输时也一样,只是多传输下更明显。需要在日志 / 排查里区分某条会话走的是哪种传输,得自行包一层 Stream 带标记(OwnedStream 可作参考),或按 ClientId 约定区分。
  21. MultiTransport 不改变「新连接顶掉同 ClientId 旧会话」的语义:一个 TCP 客户端与一个管道客户端取同一个 ClientId 会互相顶掉,且是跨传输地顶掉。多传输场景请保证 ClientId 全局唯一。

版本历史

1.2.0

本版新增 IPC 信封协议层(OpenRobot.Framework.Communication.Ipc.Envelope 命名空间)。 代码自 SmartEdgeHub 的 EdgeHub.CloudHost 抽入,接入方不必再逐项目拷贝。

  • IpcEnvelope(五字段信封 DTO)/ IpcWire(线格式常量 edge.message)/ IpcMessageContext(处理器上下文 + 回执)
  • IIpcMessageHandler(入方向扩展点)/ IIpcConnectionHook(出方向连接钩子)
  • IIpcEnvelopeSender / IpcEnvelopeSender(出方向门面:单向 + 请求-响应)
  • IIpcEnvelopeGateway / IpcEnvelopeGateway(入方向网关,独占帧类型后按 topic 二次分发)
  • 本机名改由 Attach 从 client.ClientId 取,构造函数不再吃任何配置。 原来是读配置里的 ClientId,而库解析客户端标识用三级回退(ClientId → Name → 生成的 Guid)—— 两者在 ClientId 留空时结果不同,会让信封的 source(也就是回执的 Target)填错, 触发「消息发得出去、回执转不回来」的静默故障。现在两者同源
  • 信封层用 ILogger 记日志(与库其余部分用 ICommunicationLogger 的 Ipc / Mqtt 不同)—— 这层是从宿主应用抽上来的,保留了原有的日志形状
  • 公共 API 定型(发布后即为破坏性变更,故在本版落定): IIpcEnvelopeGateway 不再继承 IIpcEnvelopeSender(去掉让 DI 环在类型上合法的写法, 注错方向从「启动时崩」变成编译错误);IpcMessageContext 构造函数由 internal 放开为 public (接入方才能单测自己的处理器,不必起真管道);IpcEnvelope 五个属性补显式 [JsonPropertyName], 把「字段名就是契约」从文档搬进代码(裸 JsonSerializer 反序列化不再静默解出全 null)
  • 未改动 IpcClient / IpcService / IpcOptions / IpcFrame 等既有类型

1.1.0

本版新增 SignalR 模块与 IPC 多传输能力。

SignalR(OpenRobot.Framework.Communication.Signalr 命名空间)

以 WebApiProxy v2 的功能面为基线重写,不迁入 V1 的 ClientSignalrProxy / ISignalrOnEvent / ProtoLogLevel / Options.HubEvents。WebApiProxy 类库本身零改动。

  • 连接控制:StartAsync(非阻塞 + 自动重连)/ ConnectAsync(单次,失败抛出)/ StopAsync(保留处理器,可重启)/ DisconnectAsync,显式状态机与 ConnectionStateChanged 事件
  • 单一重连权威:不配置 SignalR 内建 WithAutomaticReconnect,退避与尝试次数由本库循环独占决策(MaxReconnectAttempts = 0 是真的无限,退避序列用尽后复用最后一项)
  • 处理器跨重连存活:连接对象全程只构建一次,On<T1..T8> 注册的处理器不因重连丢失;返回的句柄只摘除自己那一个
  • 普通调用双形态:InvokeAsync(异常传播)/ TryInvokeAsync(Result 模式)+ 错误分类 SignalrErrorType(NotConnected / Timeout / ConnectionLost / ServerError / Unknown);用户主动取消原样抛 OperationCanceledException
  • 服务端→客户端 流:StreamAsync 返回可 Pull(ReadAllAsync)可 Push(Subscribe)的订阅,广播语义(每消费者各收全量,不是瓜分),断线自动重建,onNext 异常逐项隔离
  • 客户端→服务端 流:CreateUploadStream(逐项)/ CreateBinaryUploadStream(任意长度 byte[] 自动分片),断线后缓冲内数据继续发送
  • 空闲超时看门狗:可按「多长时间收不到业务数据」配置自动重连 / 停止 / 置失败,附可选探活钩子
  • 协议可切 MessagePack / Json,ConfigureProtocolServices 逃生舱;授权预留 AccessTokenProvider / Headers
  • 多连接:CommunicationOptions.AddSignalr(name, …),各连接独立配置、独立状态

IPC

  • 新增 MultiTransport(服务端专用):把多个监听端点合成一个 IIpcTransport,经 IpcOptions.TransportOverride 注入后,同一个 IpcService 可同时接受命名管道与 TCP 的客户端连接,两种客户端之间可互相路由 —— 会话表只有一张,SendToAsync / BroadcastAsync / RequestAsync 对它们一视同仁
  • 未改动 IpcService / IpcSession / IpcOptions:会话层只认 Stream,本就不关心流来自哪种传输
  • 客户端无需改动:每个 IpcClient 仍按 IpcOptions.Transport 选自己要用的那一种
  • 多传输下的行为约束:MaxClients 为所有传输的合计上限;同一 ClientId 跨传输会互相顶掉,需保证全局唯一;会话的 RemoteEndpoint 是监听端点的拼接串,无法反查来源传输

通用

  • AddOpenRobotCommunication 支持 AddSignalr,SignalR 连接同样按 keyed singleton 注册
  • 依赖新增 Microsoft.AspNetCore.SignalR.Client 8.0.8 与 Microsoft.AspNetCore.SignalR.Protocols.MessagePack 8.0.8(程序集 Version → 1.1.0)
  • 公共 API 不含 SignalR 客户端类型(第三方类型全部收在 Signalr/Internal/)

1.0.0(首次发布)

MQTT

  • 连接控制:StartAsync(非阻塞 + 自动重连)/ ConnectAsync / StopAsync,显式状态机与 ConnectionStateChanged 事件
  • 断线自动重连,重连后重订阅全部 filter(幂等,不重复 SUBSCRIBE)+ 按序补发(上限 / 重试次数可配)
  • 订阅以 (Topic, Handler) 批量注册(TopicFilter 去重),收包经 Channel 消费队列派发,命中全部匹配订阅(按 filter 明确度降序)
  • 处理器返回值非 null 且消息带 ResponseTopic 时自动回执(带 CorrelationData);抛异常自动回错误载荷
  • IMqttMessageInterceptor 收包前置钩子(存档 / ACK 等)
  • QoS / Retain / ResponseTopic / CorrelationData 透传;IMqttCredentialProvider 只读凭据接入

IPC

  • 一个服务端 + 多个客户端:按握手 ClientId 建会话表,定向发送 / 广播 / 请求-响应 / 单会话断开互不影响;重复 ClientId 以新连接为准
  • IIpcClient 与 IIpcSession 同构(共用 IIpcEndpoint),双向通讯天然对称
  • 单向发送走 Channel 待发队列(不等待写出);请求-响应带超时(超时返回 null,不抛)
  • 客户端断线自动重连,每次重连重发握手帧
  • 传输可切命名管道(Windows)/ UDS(Linux、macOS)/ TCP(明文),IIpcTransport 为可替换接缝
  • 处理器支持服务级(对所有会话生效)/ 会话级(优先),4 种注册重载(含强类型载荷)

通用

  • ICommunicationLogger 日志抽象(含 ConsoleCommunicationLogger / MicrosoftCommunicationLogger / NullCommunicationLogger)
  • AddOpenRobotCommunication() DI 扩展(keyed singleton,不自动连接)
  • 公共 API 不含 MQTTnet / 源项目类型
Product Compatible and additional computed target framework versions.
.NET net8.0 is compatible.  net8.0-android was computed.  net8.0-browser was computed.  net8.0-ios was computed.  net8.0-maccatalyst was computed.  net8.0-macos was computed.  net8.0-tvos was computed.  net8.0-windows was computed.  net9.0 was computed.  net9.0-android was computed.  net9.0-browser was computed.  net9.0-ios was computed.  net9.0-maccatalyst was computed.  net9.0-macos was computed.  net9.0-tvos was computed.  net9.0-windows was computed.  net10.0 was computed.  net10.0-android was computed.  net10.0-browser was computed.  net10.0-ios was computed.  net10.0-maccatalyst was computed.  net10.0-macos was computed.  net10.0-tvos was computed.  net10.0-windows was computed. 
Compatible target framework(s)
Included target framework(s) (in package)
Learn more about Target Frameworks and .NET Standard.

NuGet packages

This package is not used by any NuGet packages.

GitHub repositories

This package is not used by any popular GitHub repositories.

Version Downloads Last Updated
1.2.2 91 9/28/2026
1.2.1 96 9/22/2026
1.2.0 98 9/22/2026
1.1.0 97 9/21/2026
1.0.0 97 9/18/2026