如何为TCP Socket双向通信类实现等待响应的异步命令发送方法?
如何为TCP Socket双向通信类实现等待响应的异步命令发送方法?
你好,你的思路方向是对的,但原代码里有几个关键问题需要修正,我来一步步给你梳理清楚:
先说说你原代码的核心问题
- Promise的调用方式错误:你试图直接调用
this.responsePromise.resolve(data),但Promise实例本身是没有resolve或reject方法的——只有创建Promise时传入的executor函数里的那两个回调函数,才能触发Promise的状态变更。 - 无法处理并发请求:单个
responsePromise变量只能处理一个请求,如果同时调用两次sendCommandWithResponse,后面的请求会覆盖前面的,导致前面的请求永远无法被响应。 - 变量清空时机问题:在
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); }; }); } }
关键实现思路总结
- 保存回调而非Promise实例:我们需要把每个请求的
resolve和reject函数保存起来,而不是Promise实例本身——这是触发Promise状态变更的唯一正确方式。 - 请求注册表的选择:
- 顺序响应场景用队列(数组),简单高效;
- 乱序或带标识的场景用Map,精准匹配请求和响应。
- 超时处理:一定要加上超时逻辑,避免因为远程服务异常导致Promise一直处于pending状态。
- 并发安全:用注册表维护多个请求,解决了单个变量无法处理并发请求的问题。
使用示例
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
相关产品推荐
相关产品推荐

