Team Ai
Datasetpublic

codekingpro/portable-devtools

sourceHugging Faceupdated 5mo agoView on Hugging Face
1likes14kdownloads
server.py174 linesDownload Raw Back to asgiref
1import asyncio2import logging3import time4import traceback5 6from .compatibility import guarantee_single_callable7 8logger = logging.getLogger(__name__)9 10 11class StatelessServer:12    """13    Base server class that handles basic concepts like application instance14    creation/pooling, exception handling, and similar, for stateless protocols15    (i.e. ones without actual incoming connections to the process)16 17    Your code should override the handle() method, doing whatever it needs to,18    and calling get_or_create_application_instance with a unique `scope_id`19    and `scope` for the scope it wants to get.20 21    If an application instance is found with the same `scope_id`, you are22    given its input queue, otherwise one is made for you with the scope provided23    and you are given that fresh new input queue. Either way, you should do24    something like:25 26    input_queue = self.get_or_create_application_instance(27        "user-123456",28        {"type": "testprotocol", "user_id": "123456", "username": "andrew"},29    )30    input_queue.put_nowait(message)31 32    If you try and create an application instance and there are already33    `max_application` instances, the oldest/least recently used one will be34    reclaimed and shut down to make space.35 36    Application coroutines that error will be found periodically (every 100ms37    by default) and have their exceptions printed to the console. Override38    application_exception() if you want to do more when this happens.39 40    If you override run(), make sure you handle things like launching the41    application checker.42    """43 44    application_checker_interval = 0.145 46    def __init__(self, application, max_applications=1000):47        # Parameters48        self.application = application49        self.max_applications = max_applications50        # Initialisation51        self.application_instances = {}52 53    ### Mainloop and handling54 55    def run(self):56        """57        Runs the asyncio event loop with our handler loop.58        """59        event_loop = asyncio.get_event_loop()60        try:61            event_loop.run_until_complete(self.arun())62        except KeyboardInterrupt:63            logger.info("Exiting due to Ctrl-C/interrupt")64 65    async def arun(self):66        """67        Runs the asyncio event loop with our handler loop.68        """69 70        class Done(Exception):71            pass72 73        async def handle():74            await self.handle()75            raise Done76 77        try:78            await asyncio.gather(self.application_checker(), handle())79        except Done:80            pass81 82    async def handle(self):83        raise NotImplementedError("You must implement handle()")84 85    async def application_send(self, scope, message):86        """87        Receives outbound sends from applications and handles them.88        """89        raise NotImplementedError("You must implement application_send()")90 91    ### Application instance management92 93    def get_or_create_application_instance(self, scope_id, scope):94        """95        Creates an application instance and returns its queue.96        """97        if scope_id in self.application_instances:98            self.application_instances[scope_id]["last_used"] = time.time()99            return self.application_instances[scope_id]["input_queue"]100        # See if we need to delete an old one101        while len(self.application_instances) > self.max_applications:102            self.delete_oldest_application_instance()103        # Make an instance of the application104        input_queue = asyncio.Queue()105        application_instance = guarantee_single_callable(self.application)106        # Run it, and stash the future for later checking107        future = asyncio.ensure_future(108            application_instance(109                scope=scope,110                receive=input_queue.get,111                send=lambda message: self.application_send(scope, message),112            ),113        )114        self.application_instances[scope_id] = {115            "input_queue": input_queue,116            "future": future,117            "scope": scope,118            "last_used": time.time(),119        }120        return input_queue121 122    def delete_oldest_application_instance(self):123        """124        Finds and deletes the oldest application instance125        """126        oldest_time = min(127            details["last_used"] for details in self.application_instances.values()128        )129        for scope_id, details in self.application_instances.items():130            if details["last_used"] == oldest_time:131                self.delete_application_instance(scope_id)132                # Return to make sure we only delete one in case two have133                # the same oldest time134                return135 136    def delete_application_instance(self, scope_id):137        """138        Removes an application instance (makes sure its task is stopped,139        then removes it from the current set)140        """141        details = self.application_instances[scope_id]142        del self.application_instances[scope_id]143        if not details["future"].done():144            details["future"].cancel()145 146    async def application_checker(self):147        """148        Goes through the set of current application instance Futures and cleans up149        any that are done/prints exceptions for any that errored.150        """151        while True:152            await asyncio.sleep(self.application_checker_interval)153            for scope_id, details in list(self.application_instances.items()):154                if details["future"].done():155                    exception = details["future"].exception()156                    if exception:157                        await self.application_exception(exception, details)158                    try:159                        del self.application_instances[scope_id]160                    except KeyError:161                        # Exception handling might have already got here before us. That's fine.162                        pass163 164    async def application_exception(self, exception, application_details):165        """166        Called whenever an application coroutine has an exception.167        """168        logging.error(169            "Exception inside application: %s\n%s%s",170            exception,171            "".join(traceback.format_tb(exception.__traceback__)),172            f"  {exception}",173        )174 
codekingpro/portable-devtools · Team Ai