【问题标题】:Pattern for reading from a transport in Twisted Python从 Twisted Python 中读取传输的模式
【发布时间】:2017-06-22 12:13:50
【问题描述】:

在 Twisted Python 中,数据被写入协议的传输,但通过覆盖 dataReceived 方法接收。是否有从传输中读取的模式?这在使用 inlineCallbacks 实现状态时会很有帮助

例如:

class SomeProtocol(Protocol):
    @defer.inlineCallbacks
    def login(self):
        self.transport.write('login')
        resp = yield self.transport.read(5, timeout=1) # this doesn't exist
        if resp != 'user:':
            raise SomeException()
        self.transport.write('admin')
        resp = transport.read(9, timeout=1)
        if resp != 'password:':
            raise SomeException()
        self.transport.write('hunter2')
        # ... etc

【问题讨论】:

    标签: python twisted


    【解决方案1】:

    多年来,已经有几次尝试实现这样的 API。没有人获得任何牵引力。我想他们现在都被抛弃了。

    原则上,这并不难实现。您只是将 dataReceived 回调(一种推式 API)转换为拉式 API。

    实际上,生成的代码很脆弱,并且往往包含更多错误。

    我认为您要解决的问题是 dataReceived 是用于解析字节流的非常低级的原语。

    对此有多种可能的解决方案。您可以尝试构建一个更高级别的基于协议的工具,该工具了解您的协议的各个方面(这基本上是 Twisted 中所有协议实现所做的)。您还可以查看诸如 tubes 之类的第三方库(它为处理字节流提供了不同的抽象)。

    【讨论】:

      【解决方案2】:

      我最终维护了一个延迟列表,以便在数据到达时回调,并缓冲传入的数据,直到它满足列表中第一个延迟所需的数据长度。

      class SomeProtocol(Protocol):
      
          # initialise self.buf and self.readers in __init__
      
          def deferred_read(self, count, timeout=None):
              """Return a deferred that fires when data becomes available"""
              d = defer.Deferred()
              reader = [d, count]
              timeout_cb = None
              if timeout is not None:
                  timeout_cb = self.reactor.callLater(timeout, self.deferred_read_timeout, reader)
              reader.append(timeout_cb)
              self.readers.append(reader)
              self.check_readers()
              return d
      
          def deferred_read_timeout(self, reader):
              """Timeout this reader and check if others now match"""
              d, count, timeout_cb = reader
              self.readers.remove(reader)
              d.errback(TimeoutException()) # defined elsewhere
              self.check_readers()
      
          def check_readers(self):
              """Check if there is enough data to satisfy first reader"""
              try:
                  while 1:
                      reader = self.readers[0]
                      d, count, timeout_cb = reader
                      if len(self.buf) < count:
                          break
                      data = self.buf[:count]
                      self.buf = self.buf[count:]
                      self.readers.remove(reader)
                      try:
                          timeout_cb.cancel()
                      except: pass
                      d.callback(data)
              except IndexError: pass
      
          def dataReceived(self, data):
              self.buf += data
              self.check_readers()
      

      目前要求计数不为零。最好扩展它以支持返回当前在读取缓冲区中的任何内容,并在超时但不计数的情况下进行读取,以便在超时后返回缓冲区中的任何内容。

      【讨论】:

      • 这已经有好几年了,但我想指出,这就像用锤子试图拧开一个合适的管道。如果发送和接收之间确实有很长的时间,那么它可能是合理的,但通常当您需要同步读/写时,您不应该使用异步协议。也就是说,当返回的数据不是自识别的时,应该优先选择同步通信而不是异步通信。仅仅因为你使用了 twisted 并不意味着一切都需要基于 reactor...
      • @AlexBaum 我很好奇你将如何在 Twisted 中实现同步读取而不阻塞反应器?你有例子吗?
      • 我根本不会依赖扭曲的架构。这就是如何。这并不是说反应堆没有用,但在过去的几年里,我发现扭曲的反应堆,尽管它做了很多好事,但也做了很多坏事。我发现的两个最大弊端是它迫使您使用相当多的设置来执行单元测试,这阻止了大多数用户实际编写它们。第二个是通过使用 callLater 和 deferreds 隐藏堆栈跟踪。如果您希望我提供一个同步非扭曲 TCP 客户端的示例,我可以。
      • 问题是关于将其安装到 Twisted 中的模式,因为它是更大系统的一部分 - 我同意如果您只是执行这项任务,那么它更适合同步环境。我真的很喜欢使用trail 进行测试,尤其是能够控制时钟对我们的工作来说非常强大。这可能是适合这项工作的正确工具的问题。
      猜你喜欢
      • 2016-09-08
      • 2011-02-09
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2011-07-24
      • 1970-01-01
      • 2011-01-29
      • 2016-04-16
      相关资源
      最近更新 更多