|
| 1 | +#!/usr/bin/env python |
| 2 | +## |
| 3 | +## This file is part of the OpenSIPS Python Package |
| 4 | +## (see https://github.com/OpenSIPS/python-opensips). |
| 5 | +## |
| 6 | +## This program is free software: you can redistribute it and/or modify |
| 7 | +## it under the terms of the GNU General Public License as published by |
| 8 | +## the Free Software Foundation, either version 3 of the License, or |
| 9 | +## (at your option) any later version. |
| 10 | +## |
| 11 | +## This program is distributed in the hope that it will be useful, |
| 12 | +## but WITHOUT ANY WARRANTY; without even the implied warranty of |
| 13 | +## MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the |
| 14 | +## GNU General Public License for more details. |
| 15 | +## |
| 16 | +## You should have received a copy of the GNU General Public License |
| 17 | +## along with this program. If not, see <http://www.gnu.org/licenses/>. |
| 18 | +## |
| 19 | + |
| 20 | + |
| 21 | +""" Module that implements OpenSIPS Event behavior with asyncio """ |
| 22 | + |
| 23 | +from ..mi import OpenSIPSMIException |
| 24 | +from .json_helper import extract_json |
| 25 | +from .event import OpenSIPSEventException |
| 26 | +import asyncio |
| 27 | + |
| 28 | +class AsyncOpenSIPSEvent(): |
| 29 | + |
| 30 | + """ Asyncio implementation of the OpenSIPS Event """ |
| 31 | + |
| 32 | + def __init__(self, handler, name: str, callback, expire=None): |
| 33 | + self._handler = handler |
| 34 | + self.name = name |
| 35 | + self.callback = callback |
| 36 | + self.buf = b"" |
| 37 | + self.json_queue = [] |
| 38 | + self.retries = 0 |
| 39 | + if expire is not None: |
| 40 | + self.expire = expire |
| 41 | + self.reregister = False |
| 42 | + else: |
| 43 | + self.expire = 3600 |
| 44 | + self.reregister = True |
| 45 | + |
| 46 | + try: |
| 47 | + self.socket = self._handler.__new_socket__() |
| 48 | + self.socket.create() |
| 49 | + self._handler.events[self.name] = self |
| 50 | + self.resubscribe_task = asyncio.create_task(self.resubscribe()) |
| 51 | + loop = asyncio.get_running_loop() |
| 52 | + loop.add_reader(self.socket.sock.fileno(), self.handle, self.callback) |
| 53 | + |
| 54 | + except ValueError as e: |
| 55 | + raise OpenSIPSEventException("Invalid arguments for socket creation: {}".format(e)) |
| 56 | + |
| 57 | + def handle(self, callback): |
| 58 | + """ Handles the event callbacks """ |
| 59 | + data = self.socket.read() |
| 60 | + if not data: |
| 61 | + return |
| 62 | + |
| 63 | + self.buf += data |
| 64 | + self.json_queue, self.buf = extract_json(self.json_queue, self.buf) |
| 65 | + |
| 66 | + if not self.json_queue: |
| 67 | + self.retries += 1 |
| 68 | + |
| 69 | + if self.retries > 10: |
| 70 | + callback(None) |
| 71 | + return |
| 72 | + |
| 73 | + while self.json_queue: |
| 74 | + self.retries = 0 |
| 75 | + json_obj = self.json_queue.pop(0) |
| 76 | + callback(json_obj) |
| 77 | + |
| 78 | + async def resubscribe(self): |
| 79 | + """ Resubscribes for the event """ |
| 80 | + try: |
| 81 | + while True: |
| 82 | + try: |
| 83 | + self._handler.__mi_subscribe__(self.name, self.socket.sock_name, self.expire) |
| 84 | + except OpenSIPSEventException as e: |
| 85 | + return |
| 86 | + except OpenSIPSMIException as e: |
| 87 | + return |
| 88 | + await asyncio.sleep(self.expire - 60) |
| 89 | + except asyncio.CancelledError: |
| 90 | + pass |
| 91 | + |
| 92 | + def unsubscribe(self): |
| 93 | + """ Unsubscribes the event """ |
| 94 | + try: |
| 95 | + self._handler.__mi_unsubscribe__(self.name, self.socket.sock_name) |
| 96 | + self.stop() |
| 97 | + del self._handler.events[self.name] |
| 98 | + except OpenSIPSEventException as e: |
| 99 | + raise e |
| 100 | + except OpenSIPSMIException as e: |
| 101 | + raise e |
| 102 | + |
| 103 | + def stop(self): |
| 104 | + """ Stops the current event processing """ |
| 105 | + loop = asyncio.get_running_loop() |
| 106 | + loop.remove_reader(self.socket.sock.fileno()) |
| 107 | + self.resubscribe_task.cancel() |
| 108 | + self.socket.destroy() |
0 commit comments