-
Notifications
You must be signed in to change notification settings - Fork 5
Twisted Deferreds Example
Twisted uses an abstraction called "Deferreds" for writing non-blocking asynchronous code. They should be used whenever there is the possibility that code will block on I/O. The following is an example of how they can be used including a server that sends out chunks of a large file to clients and when the clients reply with the SHA256 hash of that chunk the server sends the next chunk. The code, by default uses deferreds, but can also be run without deferreds by changing the buildProtocol function in the Server Factory to use True instead of False: return BigFileServerProtocol(self, False) vs. return BigFileServerProtocol(self, True).
Be aware that as of right now I am not even sure if this is the correct way to use deferreds. This was produced as part of an effort to learn to use them correctly. I may have to modify this file repeatedly as I learn more.
Notably there is also a function called deferToThread which allows us to prevent blocking on CPU intensive actions which is not demonstrated here.
##Server
# Read username, output from non-empty factory, drop connections
# Use deferreds, to minimize synchronicity assumptions
from time import time
from twisted.internet import reactor, defer
from twisted.internet.protocol import Factory
from twisted.internet.protocol import Protocol
from Crypto.Hash import SHA256
def sha256(message):
return bytes(SHA256.new(message).hexdigest(), 'utf-8')
class BigFileServerProtocol(Protocol):
def __init__(self, factory, use_deferreds=True):
self.file_path = "Big File.txt"
self.current_hash = None
self.factory = factory
self.factory.protocols.append(self)
self.data_processed = 0
self.start = None
self.generator = self.file_generator()
self.defer = None
self.use_deferreds = use_deferreds
def file_generator(self, chunk_size=8192):
self.start = int(time())
with open(self.file_path, 'rb') as f:
while 42:
data = f.read(chunk_size)
if not data:
break
self.data_processed += len(data)
yield data
return None
def next_chunk(self):
return self.generator.__next__()
def connectionMade(self):
data = self.next_chunk()
self.current_hash = sha256(data)
self.transport.write(data)
def processData(self, data):
if self.use_deferreds:
d = defer.Deferred()
if data == self.current_hash:
if self.use_deferreds:
d.callback(self)
else:
data = self.next_chunk()
self.current_hash = sha256(data)
if data:
self.transport.write(data)
self.factory.update()
else:
if self.use_deferreds:
d.errback(self)
else:
self.factory.protocols.remove(self)
self.transport.loseConnection()
if self.use_deferreds:
return d
def dataReceived(self, data):
if self.use_deferreds:
def handle_error(self):
self.factory.protocols.remove(self)
self.transport.loseConnection()
def handle_success(self):
data = self.next_chunk()
self.current_hash = sha256(data)
if data:
self.transport.write(data)
self.factory.update()
d = self.processData(data)
d.addCallbacks(handle_success, handle_error)
else:
self.processData(data)
class BigFileServerFactory(Factory):
def __init__(self):
self.protocols = []
self.str_len = 0
def buildProtocol(self, address):
return BigFileServerProtocol(self, False)
def update(self):
s = ''
i = 1
for protocol in self.protocols:
if int(time() - protocol.start):
m = 1
modifier = 'B'
speed = protocol.data_processed // int(time() - protocol.start)
if speed // (1024):
m = 1024
modifier = 'KB'
if speed // (1024^2):
m = pow(1024,2)
modifier = 'MB'
if speed // pow(1024,3):
m = (1024^3)
modifier = 'GB'
speed //= m
s += "{0}: {1} {2}/s | ".format(i, speed, modifier)
else:
s += "{0}: {1} B/s | ".format(i, protocol.data_processed)
i += 1
s += ' ' * 5
print('\b' * self.str_len + s, end='')
self.str_len = len(s)
class BigFileServer(object):
def __init__(self, port):
self.port = port
def run(self):
reactor.listenTCP(self.port, BigFileServerFactory())
reactor.run()
if __name__ == '__main__':
server = BigFileServer(1024)
server.run()##Client
from os.path import isfile
from Crypto.Hash import SHA256
from twisted.internet import reactor
from twisted.internet.protocol import Protocol, ClientFactory
def sha256(message):
return bytes(SHA256.new(message).hexdigest(), 'utf-8')
class BigFileClientProtocol(Protocol):
def __init__(self, factory):
self.factory = factory
filename = "Big File"
n = ''
i = 0
while isfile(filename + n + '.txt'):
i += 1
n = " ({0})".format(i)
self.file_path = filename + n + '.txt'
def write_chunk(self, data):
h = sha256(data)
with open(self.file_path, 'ab+') as f:
f.write(data)
self.transport.write(h)
def dataReceived(self, data):
self.write_chunk(data)
class BigFileClientFactory(ClientFactory):
def __init__(self):
print("Client initiated")
def buildProtocol(self, addr):
return BigFileClientProtocol(self)
def clientConnectionFailed(self, connector, reason):
print("Connection failed.")
reactor.stop()
def clientConnectionLost(self, connector, reason):
print("Connection lost.")
reactor.stop()
class BigFileClient(object):
def __init__(self, url, port):
self.server_url = url
self.server_port = port
def run(self):
reactor.connectTCP(self.server_url, self.server_port, BigFileClientFactory())
reactor.run()
if __name__ == '__main__':
client = BigFileClient('localhost', 1024)
client.run()