|
| 1 | +import asyncio |
| 2 | +import logging |
| 3 | +from ..client import Client |
| 4 | +from ..client.abc.events import Events |
| 5 | + |
| 6 | + |
| 7 | +class Emitter(Events): |
| 8 | + |
| 9 | + _ev_handlers = dict() |
| 10 | + |
| 11 | + def __init__(self, client: Client, emitter: str, scope=None): |
| 12 | + super().__init__() |
| 13 | + self._event_id = 0 |
| 14 | + self._client = client |
| 15 | + self._thing_id = None |
| 16 | + self._scope = scope or client.get_scope() |
| 17 | + self._code = \ |
| 18 | + f'{{{emitter}}}.watch(); {{{emitter}}}.id();' |
| 19 | + client.add_event_handler(self) |
| 20 | + asyncio.ensure_future(self._watch()) |
| 21 | + |
| 22 | + def __init_subclass__(cls): |
| 23 | + cls._ev_handlers = {} |
| 24 | + |
| 25 | + for key, val in cls.__dict__.items(): |
| 26 | + if not key.startswith('__') and \ |
| 27 | + callable(val) and hasattr(val, '_ev'): |
| 28 | + cls._ev_handlers[val._ev] = val |
| 29 | + |
| 30 | + async def _watch(self): |
| 31 | + self._thing_id = await self._client.query( |
| 32 | + self._code, |
| 33 | + scope=self._scope) |
| 34 | + |
| 35 | + def on_reconnect(self): |
| 36 | + asyncio.ensure_future(self._watch()) |
| 37 | + |
| 38 | + def on_node_status(self, _status): |
| 39 | + pass |
| 40 | + |
| 41 | + def on_warning(self, warn): |
| 42 | + logging.warning(f'{warn["warn_msg"]} ({warn["warn_code"]})') |
| 43 | + |
| 44 | + def on_watch_init(self, data): |
| 45 | + pass |
| 46 | + |
| 47 | + def on_event(self, ev, *args): |
| 48 | + cls = self.__class__ |
| 49 | + fun = cls._ev_handlers.get(ev) |
| 50 | + if fun is None: |
| 51 | + logging.debug(f'no event handler for {ev} on {cls.__name__}') |
| 52 | + return |
| 53 | + fun(self, *args) |
| 54 | + |
| 55 | + def on_watch_update(self, data): |
| 56 | + thing_id = data['#'] |
| 57 | + if thing_id != self._thing_id: |
| 58 | + return |
| 59 | + |
| 60 | + event_id, jobs = data['event'], data.pop('jobs') |
| 61 | + |
| 62 | + if self._event_id > event_id: |
| 63 | + logging.warning( |
| 64 | + f'ignore event because the current event `{self._event_id}` ' |
| 65 | + f'is greather than the received event `{event_id}`') |
| 66 | + return |
| 67 | + self._event_id = event_id |
| 68 | + |
| 69 | + for job_dict in jobs: |
| 70 | + for name, job in job_dict.items(): |
| 71 | + if name == 'event': |
| 72 | + self.on_event(*job) |
| 73 | + |
| 74 | + def on_watch_delete(self, data): |
| 75 | + thing_id = data['#'] |
| 76 | + if thing_id == self._thing_id: |
| 77 | + logging.debug(f'emitter with id {thing_id} is removed') |
| 78 | + self._client.remove_event_handler(self) |
| 79 | + |
| 80 | + def on_watch_stop(self, data): |
| 81 | + thing_id = data['#'] |
| 82 | + if thing_id == self._thing_id: |
| 83 | + logging.debug(f'emitter with id {thing_id} is stopped') |
| 84 | + self._client.remove_event_handler(self) |
0 commit comments