| 28 | |
| 29 | |
| 30 | class Client(MessagingHandler): |
| 31 | def __init__(self, host, address): |
| 32 | super(Client, self).__init__() |
| 33 | self.host = host |
| 34 | self.address = address |
| 35 | self.sent = [] |
| 36 | self.pending = [] |
| 37 | self.reply_address = None |
| 38 | self.sender = None |
| 39 | self.receiver = None |
| 40 | |
| 41 | def on_start(self, event): |
| 42 | conn = event.container.connect(self.host) |
| 43 | self.sender = event.container.create_sender(conn, self.address) |
| 44 | self.receiver = event.container.create_receiver(conn, None, dynamic=True) |
| 45 | |
| 46 | def on_link_opened(self, event): |
| 47 | if event.receiver == self.receiver: |
| 48 | self.reply_address = event.link.remote_source.address |
| 49 | self.do_request() |
| 50 | |
| 51 | def on_sendable(self, event): |
| 52 | self.do_request() |
| 53 | |
| 54 | def on_message(self, event): |
| 55 | if self.sent: |
| 56 | request, future = self.sent.pop(0) |
| 57 | print("%s => %s" % (request, event.message.body)) |
| 58 | future.set_result(event.message.body) |
| 59 | self.do_request() |
| 60 | |
| 61 | def do_request(self): |
| 62 | if self.pending and self.reply_address and self.sender.credit: |
| 63 | request, future = self.pending.pop(0) |
| 64 | self.sent.append((request, future)) |
| 65 | req = Message(reply_to=self.reply_address, body=request) |
| 66 | self.sender.send(req) |
| 67 | |
| 68 | def request(self, body): |
| 69 | future = Future() |
| 70 | self.pending.append((body, future)) |
| 71 | self.do_request() |
| 72 | self.container.touch() |
| 73 | return future |
| 74 | |
| 75 | |
| 76 | class ExampleHandler(tornado.web.RequestHandler): |