codekingpro/portable-devtools
114k
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 