【问题标题】:How can I control program flow using events and promises?如何使用事件和承诺控制程序流程?
【发布时间】:2015-09-29 07:07:39
【问题描述】:

我有这样的课:

import net from 'net';
import {EventEmitter} from 'events';
import Promise from 'bluebird';

class MyClass extends EventEmitter {
    constructor(host = 'localhost', port = 10011) {
        super(EventEmitter);
        this.host = host;
        this.port = port;
        this.socket = null;
        this.connect();
    }
    connect() {
        this.socket = net.connect(this.port, this.host);
        this.socket.on('connect', this.handle.bind(this));
    }
    handle(data) {
        this.socket.on('data', data => {

        });
    }
    send(data) {
        this.socket.write(data);
    }
}

如何将send 方法变成一个promise,它从套接字的data 事件中返回一个值?服务器只在数据发送到它时才返回数据,而不是很容易被抑制的连接消息。

我尝试过类似的方法:

handle(data) {
    this.socket.on('data', data => {
        return this.socket.resolve(data);
    });
    this.socket.on('error', this.socket.reject.bind(this));
}
send(data) {
    return new Promise((resolve, reject) => {
        this.socket.resolve = resolve;
        this.socket.reject = reject;
        this.socket.write(data);
    });
}

显然这不起作用,因为resolve/reject 在并行链接和/或调用send 时会相互覆盖。

还有一个问题是同时调用send 两次,它会解决首先返回的响应。

我目前有一个使用队列和 defers 的实现,但感觉很乱,因为队列一直在被检查。

我希望能够做到以下几点:

let c = new MyClass('localhost', 10011);
c.send('foo').then(response => {
    return c.send('bar', response.param);
    //`response` should be the data returned from `this.socket.on('data')`.
}).then(response => {
    console.log(response);
}).catch(error => console.log(error));

补充一点,我对接收到的数据没有任何控制权,这意味着它不能在流之外进行修改。

编辑:所以这似乎是不可能的,因为 TCP 没有请求-响应流。如何仍然使用 Promise 来实现这一点,但使用单次执行(一次一个请求)的 Promise 链或队列。

【问题讨论】:

  • 你的意思是像双向聊天?发送一条消息,然后等到收到一条消息,就这样?
  • @thefourtheye 差不多,除了我可能需要并行调用send 并且承诺应该根据发送的内容返回正确的响应。虽然所有接收到的数据都来自一个流,所以它并不完全可追溯。
  • 我猜这里...你能不能用.add()方法在socket中设置某种观察者对象,然后从send()调用this.socket.observer.add({ reject: reject, resolve: resolve };
  • 你的意思是像Q-Connection?
  • 在您描述的方法(您尝试过的方法)中,链接时覆盖不是问题,因为您只在 then 处理程序内第二次调用 send(即,在第一个承诺已解决)。关于并行sends,无论语言/代码结构如何,这是不可能的,因为问题出在协议定义中。如果你想要一个串行的请求/响应通信协议(即没有相关的消息 id),你必须遵守规则并在发送下一个请求之前等待响应。

标签: javascript node.js tcp promise ecmascript-6


【解决方案1】:

我将问题提炼到最低限度,并使其可在浏览器中运行:

  1. Socket 类被模拟。
  2. EventEmitter 中删除了有关端口、主机和继承的信息。

该解决方案通过将新请求附加到承诺链来工作,但在任何给定时间点最多允许一个打开/未应答的请求。 .send 每次调用时都会返回一个新的 Promise,并且该类负责所有内部同步。因此,.send 可能会被多次调用,并保证请求处理的正确顺序 (FIFO)。如果没有待处理的请求,我添加的一项附加功能是修剪承诺链。


警告我完全省略了错误处理,但无论如何它都应该针对您的特定用例进行定制。


DEMO

class SocketMock {

  constructor(){
    this.connected = new Promise( (resolve, reject) => setTimeout(resolve,200) ); 
    this.listeners = {
  //  'error' : [],
    'data' : []
    }
  }

  send(data){

    console.log(`SENDING DATA: ${data}`);
    var response = `SERVER RESPONSE TO: ${data}`;
    setTimeout( () => this.listeners['data'].forEach(cb => cb(response)),               
               Math.random()*2000 + 250); 
  }

  on(event, callback){
    this.listeners[event].push(callback); 
  }

}

class SingleRequestCoordinator {

    constructor() {
        this._openRequests = 0; 
        this.socket = new SocketMock();
        this._promiseChain = this.socket
            .connected.then( () => console.log('SOCKET CONNECTED'));
      this.socket.on('data', (data) => {
        this._openRequests -= 1;
        console.log(this._openRequests);
        if(this._openRequests === 0){
          console.log('NO PENDING REQUEST --- trimming the chain');
          this._promiseChain = this.socket.connected
        }
        this._deferred.resolve(data);
      });

    }

    send(data) {
      this._openRequests += 1;
      this._promiseChain = this._promiseChain
        .then(() => {
            this._deferred = Promise.defer();
            this.socket.send(data);
            return this._deferred.promise;
        });
      return this._promiseChain;
    }
}

var sender = new SingleRequestCoordinator();

sender.send('data-1').then(data => console.log(`GOT DATA FROM SERVER --- ${data}`));
sender.send('data-2').then(data => console.log(`GOT DATA FROM SERVER --- ${data}`));
sender.send('data-3').then(data => console.log(`GOT DATA FROM SERVER --- ${data}`));

setTimeout(() => sender.send('data-4')
    .then(data => console.log(`GOT DATA FROM SERVER --- ${data}`)), 10000);

【讨论】:

    【解决方案2】:

    如果您的send() 调用相互混淆,您应该将其保存到缓存中。为确保收到的消息与发送的消息匹配,您应该为每条消息分配一些唯一的 id 到有效负载中。

    所以你的消息发送者看起来像这样

    class MyClass extends EventEmitter {
      constructor() {
        // [redacted]
        this.messages = new Map();
      }
    
      handle(data) {
        this.socket.on('data', data => {
           this.messages.get(data.id)(data);
           this.messages.delete(data.id);
        });
      }
    
      send(data) {
        return return new Promise((resolve, reject) => {
            this.messages.set(data.id, resolve);
            this.socket.write(data);
        });
      }
    }
    

    此代码对消息顺序不敏感,您将获得所需的 API。

    【讨论】:

    • 无法在脚本之外分配idid 不会从流中返回,这意味着脚本将不知道收到了哪个响应。
    • 那么,你想用来自套接字的任何下一个数据帧来解决 Promise 吗?这个解决方案似乎不可靠,但你可以做到,如果你用messages数组而不是对象
    【解决方案3】:

    socket.write(data[, encoding][, callback]) 接受回调。您可以在此回调中拒绝或解决。

    class MyClass extends EventEmitter {
      constructor(host = 'localhost', port = 10011) {
        super(EventEmitter);
        this.host = host;
        this.port = port;
        this.socket = null;
        this.requests = null;
        this.connect();
      }
      connect() {
        this.socket = net.connect(this.port, this.host);
        this.socket.on('connect', () => {
          this.requests = [];
          this.socket.on('data', this.handle.bind(this));
          this.socket.on('error', this.error.bind(this));
        });
      }
      handle(data) {
        var [request, resolve, reject] = this.requests.pop();
        // I'm not sure what will happen with the destructuring if requests is empty
        if(resolve) {
          resolve(data);
        }
      }
      error(error) {
        var [request, resolve, reject] = this.requests.pop();
        if(reject) {
          reject(error);
        }
      }
      send(data) {
        return new Promise((resolve, reject) => {
          if(this.requests === null) {
            return reject('Not connected');
          }
          this.requests.push([data, resolve, reject]);
          this.socket.write(data);
        });
      }
    }
    

    未经测试,因此不确定方法签名,但这是基本思想。

    这假设每个请求都会有一个 handleerror 事件。

    我越想越觉得如果没有应用程序数据中的其他信息(例如与请求响应匹配的数据包编号),这似乎是不可能的。

    它现在的实现方式(以及它在您的问题中的方式),甚至不确定一个答案是否与一个 handle 事件完全匹配。

    【讨论】:

    • 回调只是检查数据是否发送,我想在收到数据时解决promise。唯一收到的数据来自data 事件。
    • 那么你的 Promise 就不走运了。关键是只运行一次。承诺是不管它已经发生还是将来可能发生……
    • 但是promise只会被解析一次。基本上我希望流程是send -> create promise -> return promise -> receive data -> resolve promise
    • 是的,就像我说的,它只会解决一次,这就是 Promise 的意义所在。您的用例不适合使用 Promises。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2017-12-14
    • 1970-01-01
    • 2019-11-06
    • 1970-01-01
    • 1970-01-01
    • 2014-08-25
    • 1970-01-01
    相关资源
    最近更新 更多