From 0a2cefe556e375f810ebeeb431af609f0c8213b7 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 25 Dec 2024 14:08:28 +0000 Subject: [PATCH] Fix deadlock --- trustgraph-flow/trustgraph/gateway/mux.py | 13 ++++++++++++- 1 file changed, 12 insertions(+), 1 deletion(-) diff --git a/trustgraph-flow/trustgraph/gateway/mux.py b/trustgraph-flow/trustgraph/gateway/mux.py index a35ec01b..ae699ae6 100644 --- a/trustgraph-flow/trustgraph/gateway/mux.py +++ b/trustgraph-flow/trustgraph/gateway/mux.py @@ -35,6 +35,7 @@ class MuxEndpoint(SocketEndpoint): async def maybe_tidy_workers(self, workers): while True: + try: await asyncio.wait_for( @@ -44,7 +45,8 @@ class MuxEndpoint(SocketEndpoint): # worker[0] now stopped # FIXME: Delete reference??? - del workers[0] + + workers.pop(0) if len(workers) == 0: break @@ -57,6 +59,10 @@ class MuxEndpoint(SocketEndpoint): async def start_request_task(self, ws, id, svc, request, workers): + if svc not in self.services: + await ws.send_json({"id": id, "error": "Service not recognised"}) + return + requestor = self.services[svc] async def responder(resp, fin): @@ -68,8 +74,13 @@ class MuxEndpoint(SocketEndpoint): # Wait for outstanding requests to go below MAX_OUTSTANDING_REQUESTS while len(workers) > MAX_OUTSTANDING_REQUESTS: + + # Fixes deadlock + # FIXME: Put it in its own loop await asyncio.sleep(START_REQUEST_WAIT) + await self.maybe_tidy_workers(workers) + worker = asyncio.create_task( requestor.process(request, responder) )