Files
blog-press/docs/Web/Network/NetCore.md

12 KiB
Raw Blame History

title, date
title date
.Net Core实现 2026-05-27

一、.Net Core实现

1.1 SuperSocket

SuperSocket是一个轻量级, 跨平台而且可扩展的 .Net/Mono Socket 服务器程序框架。可以轻松构建TCP、UDP、WebSocket服务器。

1.2 安装依赖

NuGut安装SuperSocket、SuperSocket.WebSocket和SuperSocket.WebSocket.Server 2.0及以上版本。

1.3 配置文件

appsettings.json

{
  "serverOptions": {
    "TcpServer": {
      "name": "TcpServer",
      "listeners": [
        {
          "ip": "Any",
          "port": 4040
        }
      ]
    },
    "WebSocketServer": {
      "name": "WebSocket",
      "listeners": [
        {
          "ip": "Any",
          "port": 5050
        }
      ]
    }
  },
  "Mqtt": {
    "Host": "127.0.0.1",
    "Port": 1883,
    "ClientId": "MqttClient",
    "Topics": [
      "test/topic1",
      "test/topic2"
    ]
  }
}

1.4 主程序

program.cs

var host = Host.CreateDefaultBuilder(args)
    .ConfigureServices((context, services) =>
    {
        services.AddSingleton<IPackageHandler<ProtocolFrame>, TcpPackageHandler>();
        services.AddSingleton<WebSocketMessageHandler>();
        services.AddSingleton<WebSocketConnectionManager>();

        // 1. 注入 MQTT 配置
        services.Configure<MqttSettings>(context.Configuration.GetSection("Mqtt"));
        // 2. 注册 MQTT 客户端
        services.AddSingleton<IMqttClient>(serviceProvider => new MqttClientFactory().CreateMqttClient());
        // 3. 注册 MQTT 后台服务
        services.AddHostedService<MqttService>();
    })
    .AsMultipleServerHostBuilder()
    .AddServer<TcpService, ProtocolFrame, BinaryPipelineFilter>(builder =>
    {
        builder.ConfigureServerOptions((ctx, config) => config.GetSection("TcpServer"));
    })
    .AddWebSocketServer(builder =>
    {
        builder
            .UseSessionHandler(
                session =>
                {
                    if (session is WebSocketSession wsSession)
                    {
                        var connectionManager = session.Server.ServiceProvider.GetRequiredService<WebSocketConnectionManager>();
                        connectionManager.Add(wsSession);
                    }
                    return ValueTask.CompletedTask;
                },
                (session, reason) =>
                {
                    if (session is WebSocketSession wsSession)
                    {
                        var connectionManager = session.Server.ServiceProvider.GetRequiredService<WebSocketConnectionManager>();
                        connectionManager.Remove(wsSession);
                    }
                    return ValueTask.CompletedTask;
                }
            )
            .UseWebSocketMessageHandler(async (session, package) =>
            {
                using var scope = session.Server.ServiceProvider.CreateScope();
                var handler = scope.ServiceProvider.GetRequiredService<WebSocketMessageHandler>();

                await handler.HandleAsync(session, package);
            })
            .ConfigureServerOptions((ctx, config) => config.GetSection("WebSocketServer"));
    })
    .ConfigureLogging(logging => logging.ClearProviders())
    .UseNLog()
    .Build();

await host.RunAsync();

  通过.AsMultipleServerHostBuilder()可以构造多服务器实例。

二、 TCP服务器

2.1 协议数据

/// <summary>
/// 协议数据结构 示例65 6D 00 05 00 01 68 65 6C 6C 6F
/// </summary>
/// <param name="Magic">帧头标识</param>
/// <param name="Length">数据长度</param>
/// <param name="Type">数据类型</param>
/// <param name="Payload">数据内容</param>
/// 
record ProtocolFrame(
    ushort Magic,
    ushort Length,
    ushort Type,
    byte[] Payload
);

  可以根据实际情况自定义消息格式。

2.2 协议解析

class BinaryPipelineFilter : FixedHeaderPipelineFilter<ProtocolFrame>
{
    /// <summary>
    /// 固定头长度
    /// </summary>
    public BinaryPipelineFilter() : base(6)
    {
    }

    /// <summary>
    /// 从包头中解析出 Body 长度
    /// </summary>
    /// <param name="buffer">字节流</param>
    /// <returns>Body 长度</returns>
    protected override int GetBodyLengthFromHeader(ref ReadOnlySequence<byte> buffer)
    {
        var reader = new SequenceReader<byte>(buffer);

        // 读取前2字节 → magic
        reader.TryReadBigEndian(out ushort magic);

        // 再读2字节 → Length
        reader.TryReadBigEndian(out ushort length);

        // 再读2字节 → Type
        reader.TryReadBigEndian(out ushort type);

        // 校验帧头
        if (magic != 0x656D)
        {
            throw new Exception("非法帧头");
        }

        // 限制长度
        if (length == 0 || length > 8192)
        {
            throw new Exception("非法长度");
        }

        // 返回 Payload 长度
        return length;
    }

    /// <summary>
    /// 把完整字节包 → 转换成 ProtocolFrame
    /// </summary>
    /// <param name="buffer">字节流</param>
    /// <returns>ProtocolFrame</returns>
    protected override ProtocolFrame DecodePackage(ref ReadOnlySequence<byte> buffer)
    {
        var reader = new SequenceReader<byte>(buffer);

        // 读取前2字节 → magic
        reader.TryReadBigEndian(out ushort magic);

        // 再读2字节 → Length
        reader.TryReadBigEndian(out ushort length);

        // 再读2字节 → Type
        reader.TryReadBigEndian(out ushort type);

        var payload = buffer.Slice(6, length).ToArray();

        // 构造ProtocolFrame
        return new ProtocolFrame(magic, length, type, payload);
    }
}

根据协议数据ProtocolFrame来解包。

2.3 消息处理

class TcpPackageHandler(ILogger<TcpPackageHandler> logger) : IPackageHandler<ProtocolFrame>
{
    private readonly ILogger<TcpPackageHandler> _logger = logger;

    public async ValueTask Handle(IAppSession session, ProtocolFrame package, CancellationToken cancellationToken)
    {
        _logger.LogInformation($"Magic={package.Magic:X4}, Type={package.Type}, Len={package.Length}");
        await session.SendAsync(package.Payload, cancellationToken);
    }
}

2.4 TCPService

class TcpService(IServiceProvider serviceProvider, IOptions<ServerOptions> serverOptions, ILogger<TcpService> logger) : SuperSocketService<ProtocolFrame>(serviceProvider, serverOptions)
{
    private readonly ILogger<TcpService> _logger = logger;

    private readonly ConcurrentDictionary<string, string> _sessionDevices = new();

    public void BindDevice(string sessionId, string deviceNum)
    {
        _sessionDevices.AddOrUpdate(sessionId, deviceNum, (k, v) => deviceNum);
    }

    protected override async ValueTask OnSessionConnectedAsync(IAppSession session)
    {
        _logger.LogInformation("TCP Session 连上: {SessionID}", session.SessionID);
        await base.OnSessionConnectedAsync(session);
    }

    protected override async ValueTask OnSessionClosedAsync(IAppSession session, CloseEventArgs e)
    {
        _logger.LogInformation("TCP Session 断开: {SessionID}, Reason={Reason}", session.SessionID, e.Reason);

        if (_sessionDevices.TryRemove(session.SessionID, out var deviceNum))
        {
            _logger.LogInformation("设备 {DeviceNum} 标记离线", deviceNum);
        }

        await base.OnSessionClosedAsync(session, e);
    }
}

三、Websocket服务器

3.1 消息处理

class WebSocketMessageHandler(ILogger<WebSocketMessageHandler> logger)
{
    private readonly ILogger<WebSocketMessageHandler> _logger = logger;

    public async ValueTask HandleAsync(WebSocketSession session, WebSocketPackage package)
    {
        _logger.LogInformation($"[WebSocket] {package.Message}");
        await session.SendAsync("ok");
    }
}

3.2 连接管理

public class WebSocketConnectionManager(ILogger<WebSocketConnectionManager> logger)
{
    private readonly ConcurrentDictionary<string, WebSocketSession> _sessions = new();
    private readonly ILogger<WebSocketConnectionManager> _logger = logger;

    public void Add(WebSocketSession session) => _sessions.TryAdd(session.SessionID, session);

    public void Remove(WebSocketSession session) => _sessions.TryRemove(session.SessionID, out _);

    public async Task BroadcastAsync(string json)
    {
        var payload = new ReadOnlyMemory<byte>(System.Text.Encoding.UTF8.GetBytes(json));

        var tasks = _sessions.Values
            .Where(s => s.State == SessionState.Connected)
            .Select(async s =>
            {
                try
                {
                    await s.SendAsync(json);
                }
                catch
                {
                    Remove(s);
                }
            });

        await Task.WhenAll(tasks);
    }
}

四、 MQTT服务器

4.1 MqttService

class MqttService(IMqttClient mqttClient, IOptions<MqttSettings> settings, ILogger<TcpPackageHandler> logger) : BackgroundService
{
    private readonly IMqttClient _mqttClient = mqttClient;
    private readonly MqttSettings _settings = settings.Value;
    private readonly ILogger<TcpPackageHandler> _logger = logger;
    private MqttClientOptions? _mqttOptions;

    protected override async Task ExecuteAsync(CancellationToken stoppingToken)
    {
        // 构建连接配置
        _mqttOptions = new MqttClientOptionsBuilder()
            .WithClientId($"{_settings.ClientId}_{Guid.NewGuid():N}")
            .WithTcpServer(_settings.Host, _settings.Port)
            .WithCleanSession()
            .Build();

        _mqttClient.ConnectedAsync += OnConnectedAsync;
        _mqttClient.DisconnectedAsync += OnDisconnectedAsync;
        _mqttClient.ApplicationMessageReceivedAsync += HandleMessage;

        _logger.LogInformation("[MQTT] 正在连接到服务器 {Host}:{Port}...", _settings.Host, _settings.Port);

        await _mqttClient.ConnectAsync(_mqttOptions, stoppingToken);

        // 等待程序停止
        await Task.Delay(Timeout.Infinite, stoppingToken);
    }

    public async Task OnConnectedAsync(MqttClientConnectedEventArgs arg)
    {
        _logger.LogInformation("[MQTT] 已成功连接到服务器");

        // 连接成功后订阅主题
        foreach (var topic in _settings.Topics)
        {
            await _mqttClient.SubscribeAsync(topic, MqttQualityOfServiceLevel.AtLeastOnce);
            _logger.LogInformation($"[MQTT] 已订阅主题:{topic}");
        }

        await Task.CompletedTask;
    }

    public async Task OnDisconnectedAsync(MqttClientDisconnectedEventArgs arg)
    {
        _logger.LogError("[MQTT] 连接失败,原因:{Reason}", arg.Reason);

        await Task.CompletedTask;
    }

    public async Task HandleMessage(MqttApplicationMessageReceivedEventArgs arg)
    {
        var topic = arg.ApplicationMessage.Topic;
        var payload = Encoding.UTF8.GetString(arg.ApplicationMessage.Payload);

        _logger.LogInformation($"\n[MQTT] 收到消息\n主题{topic}\n内容{payload}\n");

        await Task.CompletedTask;
    }

    public async Task<bool> PublishAsync(string topic, string payload, CancellationToken cancellationToken = default)
    {
        try
        {
            var message = new MqttApplicationMessageBuilder()
                .WithTopic(topic)
                .WithPayload(payload)
                .WithQualityOfServiceLevel(MqttQualityOfServiceLevel.AtLeastOnce)
                .Build();

            await _mqttClient.PublishAsync(message, cancellationToken);
            return true;
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "[MQTT] 发布消息到主题 {Topic} 失败", topic);
            return false;
        }
    }

    public override async Task StopAsync(CancellationToken stoppingToken)
    {
        await _mqttClient.DisconnectAsync();
        await base.StopAsync(stoppingToken);
    }
}

  需要实现后台服务接口一直运行。