| 7 | |
| 8 | |
| 9 | class Reactor(threading.Thread): |
| 10 | def __init__(self, driver): |
| 11 | super().__init__() |
| 12 | self.driver = driver |
| 13 | self.loop = asyncio.new_event_loop() |
| 14 | self.lock = threading.Lock() |
| 15 | self.event = threading.Event() |
| 16 | self.daemon = True |
| 17 | self.handlers = {} |
| 18 | |
| 19 | def add_event_handler(self, method_name, callback): |
| 20 | """ |
| 21 | Parameters |
| 22 | ---------- |
| 23 | event_name: str |
| 24 | example "Network.responseReceived" |
| 25 | |
| 26 | callback: callable |
| 27 | callable which accepts 1 parameter: the message object dictionary |
| 28 | """ |
| 29 | with self.lock: |
| 30 | self.handlers[method_name.lower()] = callback |
| 31 | |
| 32 | @property |
| 33 | def running(self): |
| 34 | return not self.event.is_set() |
| 35 | |
| 36 | def run(self): |
| 37 | try: |
| 38 | asyncio.set_event_loop(self.loop) |
| 39 | self.loop.run_until_complete(self.listen()) |
| 40 | except Exception as e: |
| 41 | logger.warning("Reactor.run() => %s", e) |
| 42 | |
| 43 | async def _wait_service_started(self): |
| 44 | while True: |
| 45 | with self.lock: |
| 46 | if ( |
| 47 | getattr(self.driver, "service", None) |
| 48 | and getattr(self.driver.service, "process", None) |
| 49 | and self.driver.service.process.poll() |
| 50 | ): |
| 51 | await asyncio.sleep(self.driver._delay or 0.25) |
| 52 | else: |
| 53 | break |
| 54 | |
| 55 | async def listen(self): |
| 56 | while self.running: |
| 57 | await self._wait_service_started() |
| 58 | await asyncio.sleep(1) |
| 59 | try: |
| 60 | with self.lock: |
| 61 | log_entries = self.driver.get_log("performance") |
| 62 | for entry in log_entries: |
| 63 | try: |
| 64 | obj_serialized: str = entry.get("message") |
| 65 | obj = json.loads(obj_serialized) |
| 66 | message = obj.get("message") |
no outgoing calls
no test coverage detected
searching dependent graphs…