From dcd898eb52e4707c8eaeecd0ab78268125da8608 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Thu, 1 May 2025 17:26:23 +0100 Subject: [PATCH] Hacking --- .../trustgraph/gateway/dispatch/requestor.py | 6 +- .../trustgraph/gateway/dispatch/streamer.py | 99 +++++++++++++++++++ .../trustgraph/gateway/endpoint/flows.py | 2 +- 3 files changed, 105 insertions(+), 2 deletions(-) create mode 100644 trustgraph-flow/trustgraph/gateway/dispatch/streamer.py diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/requestor.py b/trustgraph-flow/trustgraph/gateway/dispatch/requestor.py index ce5b8dc7..226fc9a9 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/requestor.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/requestor.py @@ -33,13 +33,17 @@ class ServiceRequestor: self.timeout = timeout + self.running = True + async def start(self): await self.pub.start() await self.sub.start() + self.running = True async def stop(self): await self.pub.stop() await self.sub.stop() + self.running = False def to_request(self, request): raise RuntimeError("Not defined") @@ -57,7 +61,7 @@ class ServiceRequestor: await self.pub.send(id, self.to_request(request)) - while True: + while self.running: try: resp = await asyncio.wait_for( diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/streamer.py b/trustgraph-flow/trustgraph/gateway/dispatch/streamer.py new file mode 100644 index 00000000..60d4aff2 --- /dev/null +++ b/trustgraph-flow/trustgraph/gateway/dispatch/streamer.py @@ -0,0 +1,99 @@ + +import asyncio +import uuid +import logging + +from ... base import Publisher +from ... base import Subscriber + +logger = logging.getLogger("requestor") +logger.setLevel(logging.INFO) + +class ServiceRequestor: + + def __init__( + self, + pulsar_client, + queue, schema, + handler, + subscription="api-gateway", consumer_name="api-gateway", + timeout=600, + ): + + self.sub = Subscriber( + pulsar_client, queue, + subscription, consumer_name, + schema + ) + + self.timeout = timeout + + self.running = True + + self.receiver = handler + + async def start(self): + await self.sub.start() + self.streamer = asyncio.create_task(self.stream()) + sub.start() + self.running = True + + async def stop(self): + await self.sub.stop() + self.running = False + + def from_inbound(self, response): + raise RuntimeError("Not defined") + + async def stream(self): + + id = str(uuid.uuid4()) + + try: + + q = await self.sub.subscribe(id) + + while self.running: + + try: + resp = await asyncio.wait_for( + q.get(), timeout=self.timeout + ) + except Exception as e: + raise RuntimeError("Timeout") + + if resp.error: + err = { "error": { + "type": resp.error.type, + "message": resp.error.message, + } } + + fin = False + + await self.receiver(err, fin) + + else: + + resp, fin = self.from_inbound(resp) + + print(resp, fin) + + await self.receiver(resp, fin) + + if fin: break + + except Exception as e: + + logging.error(f"Exception: {e}") + + err = { "error": { + "type": "gateway-error", + "message": str(e), + } } + if responder: + await responder(err, True) + return err + + finally: + await self.sub.unsubscribe(id) + diff --git a/trustgraph-flow/trustgraph/gateway/endpoint/flows.py b/trustgraph-flow/trustgraph/gateway/endpoint/flows.py index 89ef76ed..34396518 100644 --- a/trustgraph-flow/trustgraph/gateway/endpoint/flows.py +++ b/trustgraph-flow/trustgraph/gateway/endpoint/flows.py @@ -99,7 +99,7 @@ class FlowEndpointManager: await self.services[k].stop() del self.services[k] - self.services[k] = steamer( + self.services[k] = streamer( # pulsar_client=self.pulsar_client, # timeout = self.timeout, # input_queue = intf[api_kind],