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

如何为TCP Socket双向通信类实现等待响应的异步命令发送方法?

如何为TCP Socket双向通信类实现等待响应的异步命令发送方法?

你好,你的思路方向是对的,但原代码里有几个关键问题需要修正,我来一步步给你梳理清楚:

先说说你原代码的核心问题

  1. Promise的调用方式错误:你试图直接调用this.responsePromise.resolve(data),但Promise实例本身是没有resolve或reject方法的——只有创建Promise时传入的executor函数里的那两个回调函数,才能触发Promise的状态变更。
  2. 无法处理并发请求:单个responsePromise变量只能处理一个请求,如果同时调用两次sendCommandWithResponse,后面的请求会覆盖前面的,导致前面的请求永远无法被响应。
  3. 变量清空时机问题:在processMessage里直接把this.responsePromise设为undefined,虽然这里不是最严重的问题,但逻辑上也不够严谨。

正确的实现方案

我们需要维护一个请求注册表,用来保存每个异步请求对应的resolve/reject回调,这样既能正确触发Promise状态,又能处理并发请求。下面分两种场景给出实现:

场景1:远程服务按请求顺序返回响应(发一个等一个,无乱序)

这种情况用**队列(数组)**来保存请求回调即可:

import { Socket } from 'net';

class RigControl {
  private socket: Socket = new Socket();
  private host: string = "localhost";
  private port: number = 4532;
  // 保存待处理请求的回调队列
  private pendingRequests: Array<{ resolve: (value: string) => void; reject: (reason?: any) => void }> = [];

  connect(options?: { host?: string, port?: number }): Promise<void> {
    this.host = options?.host ?? this.host;
    this.port = options?.port ?? this.port;
    this.socket.on('data', this.processMessage.bind(this));
    
    return new Promise((resolve, reject) => {
      this.socket.connect(this.port, this.host, () => resolve());
      this.socket.on('error', (err) => reject(err));
    });
  }

  processMessage(data: Buffer | string) {
    // 先做你原本的无条件处理(比如日志、状态通知解析等)
    console.log("Received raw data:", data);
    
    // 把响应转成字符串(根据你的实际编码调整,这里用ascii)
    const response = data.toString('ascii').trim();

    // 检查是否有待处理的请求,有则触发第一个请求的回调
    if (this.pendingRequests.length > 0) {
      const { resolve } = this.pendingRequests.shift()!;
      resolve(response);
    }
  }

  sendCommand(command: string) {
    this.socket.write(command);
  }

  // 新增的异步命令方法,返回Promise<string>
  async sendCommandWithResponse(command: string): Promise<string> {
    return new Promise((resolve, reject) => {
      // 将当前请求的回调加入队列
      this.pendingRequests.push({ resolve, reject });
      // 发送命令
      this.sendCommand(command);

      // 可选:添加超时处理,避免无限等待
      const timeoutTimer = setTimeout(() => {
        // 从队列中移除当前请求(如果还未被处理)
        const requestIndex = this.pendingRequests.findIndex(req => req.resolve === resolve);
        if (requestIndex !== -1) {
          this.pendingRequests.splice(requestIndex, 1);
        }
        reject(new Error(`Command "${command}" timed out after 5s`));
      }, 5000); // 5秒超时,可根据需求调整

      // 给resolve包装一下,触发时清除超时定时器
      const originalResolve = resolve;
      resolve = (value: string) => {
        clearTimeout(timeoutTimer);
        originalResolve(value);
      };
    });
  }
}

场景2:远程服务响应可能乱序,或响应带请求标识

如果你的命令和响应有一一对应的唯一标识(比如命令带ID:"CMD_123 SET_VOLUME 5",响应返回"RESP_123 OK"),那用Map来保存请求回调更可靠,能精准匹配请求和响应:

// 只修改关键部分,其他代码和上面一致
class RigControl {
  // 用请求ID作为key,保存对应的回调
  private pendingRequests: Map<string, { resolve: (value: string) => void; reject: (reason?: any) => void }> = new Map();
  private nextRequestId = 1; // 自增ID生成器

  processMessage(data: Buffer | string) {
    // 原本的无条件处理
    console.log("Received raw data:", data);
    const responseStr = data.toString('ascii').trim();

    // 解析响应中的请求ID(假设响应格式是"RESP:<ID>:<CONTENT>")
    const respMatch = responseStr.match(/RESP:(\d+):(.+)/);
    if (respMatch) {
      const reqId = respMatch[1];
      const responseContent = respMatch[2];
      // 查找对应的请求回调
      const request = this.pendingRequests.get(reqId);
      if (request) {
        request.resolve(responseContent);
        this.pendingRequests.delete(reqId);
      }
    }
  }

  async sendCommandWithResponse(command: string): Promise<string> {
    const reqId = this.nextRequestId.toString();
    this.nextRequestId++; // 自增ID,避免重复

    return new Promise((resolve, reject) => {
      // 保存回调到Map
      this.pendingRequests.set(reqId, { resolve, reject });
      // 发送带ID的命令(格式根据你的协议定义调整)
      this.sendCommand(`CMD:${reqId}:${command}`);

      // 超时处理
      const timeoutTimer = setTimeout(() => {
        if (this.pendingRequests.has(reqId)) {
          this.pendingRequests.delete(reqId);
          reject(new Error(`Command "${command}" (ID: ${reqId}) timed out`));
        }
      }, 5000);

      // 包装resolve,触发时清除超时
      const originalResolve = resolve;
      resolve = (value) => {
        clearTimeout(timeoutTimer);
        originalResolve(value);
      };
    });
  }
}

关键实现思路总结

  1. 保存回调而非Promise实例:我们需要把每个请求的resolve和reject函数保存起来,而不是Promise实例本身——这是触发Promise状态变更的唯一正确方式。
  2. 请求注册表的选择:
    • 顺序响应场景用队列(数组),简单高效;
    • 乱序或带标识的场景用Map,精准匹配请求和响应。
  3. 超时处理:一定要加上超时逻辑,避免因为远程服务异常导致Promise一直处于pending状态。
  4. 并发安全:用注册表维护多个请求,解决了单个变量无法处理并发请求的问题。

使用示例

async function runDemo() {
  const rig = new RigControl();
  try {
    await rig.connect({ host: "localhost", port: 4532 });
    
    // 发送命令并等待响应
    const volumeResp = await rig.sendCommandWithResponse("SET_VOLUME 7");
    console.log("音量设置响应:", volumeResp);

    const statusResp = await rig.sendCommandWithResponse("GET_STATUS");
    console.log("设备状态响应:", statusResp);
  } catch (err) {
    console.error("操作失败:", err);
  } finally {
    rig.socket.end(); // 记得关闭连接
  }
}

runDemo();

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 13:25:28