You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何向Kestrel中已连接的指定Socket发送消息?

解决方案

首先纠正你代码里的几个基础错误:

  • 接口属性拼写错误:[HttPost] 应该改为 [HttpPost]
  • 发送消息时的编码错误:socket.Send(Encoding.UTF8.GetString(msg)) 应该用 GetBytes 而非 GetString,因为Send方法需要字节数组

接下来是核心的连接管理和消息发送实现:

1. 实现TCP连接的Socket管理

Kestrel的ConnectionHandler会处理每个新接入的TCP连接,我们需要在自定义的MyTcpHandler中维护已连接Socket的映射关系,使用线程安全的集合来应对并发场景:

using System.Collections.Concurrent;
using System.Net;
using System.Net.Sockets;
using Microsoft.AspNetCore.Connections;

public class MyTcpHandler : ConnectionHandler
{
    // 线程安全的集合,存储IPEndpoint与对应Socket的映射
    public static ConcurrentDictionary<IPEndPoint, Socket> ConnectedSockets { get; } = new();

    public override async Task OnConnectedAsync(ConnectionContext context)
    {
        // 获取当前连接的Socket对象
        var socket = context.Transport.GetType().GetProperty("Socket")?.GetValue(context.Transport) as Socket;
        if (socket == null)
        {
            await context.DisposeAsync();
            return;
        }

        // 获取远程设备的IPEndpoint
        var remoteEndPoint = socket.RemoteEndPoint as IPEndPoint;
        if (remoteEndPoint != null)
        {
            // 将Socket加入映射集合
            ConnectedSockets.TryAdd(remoteEndPoint, socket);
        }

        try
        {
            // 保持连接存活(Kestrel需要持续读取数据,否则会断开连接)
            await context.Transport.Input.CopyToAsync(context.Transport.Output);
        }
        finally
        {
            // 连接断开时从集合中移除
            if (remoteEndPoint != null)
            {
                ConnectedSockets.TryRemove(remoteEndPoint, out _);
            }
            await context.DisposeAsync();
        }
    }
}

2. 完善HttpPost接口的实现

在你的Controller中,通过MyTcpHandler.ConnectedSockets集合获取对应Socket,并完成消息发送,同时处理异常情况:

using System.Net;
using System.Net.Sockets;
using System.Text;
using Microsoft.AspNetCore.Mvc;

[ApiController]
[Route("[controller]")]
public class DevicesController : ControllerBase
{
    // 假设这是你维护的设备ID到IPEndpoint的映射表
    private readonly Dictionary<int, IPEndPoint> _ipEndPoints = new()
    {
        { 1, new IPEndPoint(IPAddress.Parse("192.168.1.100"), 1234) },
        { 2, new IPEndPoint(IPAddress.Parse("192.168.1.101"), 1234) }
    };

    [HttpPost("devices/{id}")]
    public async Task<IActionResult> Send(int id, [FromBody] string msg)
    {
        // 验证设备ID是否存在
        if (!_ipEndPoints.TryGetValue(id, out var ipEndPoint))
        {
            return NotFound($"设备ID {id} 不存在");
        }

        // 从映射集合中获取对应的Socket
        if (!MyTcpHandler.ConnectedSockets.TryGetValue(ipEndPoint, out var socket))
        {
            return BadRequest($"设备 {ipEndPoint} 未连接");
        }

        try
        {
            // 将消息转为字节数组并发送
            var buffer = Encoding.UTF8.GetBytes(msg);
            await socket.SendAsync(buffer, SocketFlags.None);
            return Ok("消息发送成功");
        }
        catch (SocketException ex)
        {
            // 发送失败时移除无效连接
            MyTcpHandler.ConnectedSockets.TryRemove(ipEndPoint, out _);
            return StatusCode(500, $"消息发送失败:{ex.Message}");
        }
    }
}

关键说明

  • 线程安全集合:使用ConcurrentDictionary确保多线程场景下(同时有TCP连接建立/断开、HTTP请求调用)的操作安全
  • 连接存活处理:在OnConnectedAsync中执行Input.CopyToAsync(Output)是为了让Kestrel保持连接,否则Kestrel会因为没有数据交互而主动断开TCP连接
  • 异常处理:发送消息时捕获SocketException,并移除无效连接,避免后续请求使用已断开的Socket

内容的提问来源于stack exchange,提问作者Vincent

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.30 01:13:38