| 24 | |
| 25 | |
| 26 | class TornadoLoopHandler: |
| 27 | |
| 28 | def __init__(self, loop=None, handler_base=None): |
| 29 | self.loop = loop or tornado.ioloop.IOLoop.instance() |
| 30 | self.io = handler_base or IOHandler() |
| 31 | self.count = 0 |
| 32 | |
| 33 | def on_reactor_init(self, event): |
| 34 | self.reactor = event.reactor |
| 35 | |
| 36 | def on_reactor_quiesced(self, event): |
| 37 | event.reactor.yield_() |
| 38 | |
| 39 | def on_unhandled(self, name, event): |
| 40 | event.dispatch(self.io) |
| 41 | |
| 42 | def _events(self, sel): |
| 43 | events = self.loop.ERROR |
| 44 | if sel.reading: |
| 45 | events |= self.loop.READ |
| 46 | if sel.writing: |
| 47 | events |= self.loop.WRITE |
| 48 | return events |
| 49 | |
| 50 | def _schedule(self, sel): |
| 51 | if sel.deadline: |
| 52 | self.loop.add_timeout(sel.deadline, lambda: self._expired(sel)) |
| 53 | |
| 54 | def _expired(self, sel): |
| 55 | self.reactor.mark() |
| 56 | sel.expired() |
| 57 | self._process() |
| 58 | |
| 59 | def _process(self): |
| 60 | self.reactor.process() |
| 61 | if not self.reactor.quiesced: |
| 62 | self.loop.add_callback(self._process) |
| 63 | |
| 64 | def _callback(self, sel, events): |
| 65 | if self.loop.READ & events: |
| 66 | sel.readable() |
| 67 | if self.loop.WRITE & events: |
| 68 | sel.writable() |
| 69 | self._process() |
| 70 | |
| 71 | def on_selectable_init(self, event): |
| 72 | sel = event.context |
| 73 | if sel.fileno() >= 0: |
| 74 | self.loop.add_handler(sel.fileno(), lambda fd, events: self._callback(sel, events), self._events(sel)) |
| 75 | self._schedule(sel) |
| 76 | self.count += 1 |
| 77 | |
| 78 | def on_selectable_updated(self, event): |
| 79 | sel = event.context |
| 80 | if sel.fileno() >= 0: |
| 81 | self.loop.update_handler(sel.fileno(), self._events(sel)) |
| 82 | self._schedule(sel) |
| 83 | |