diff --git a/trustgraph-flow/trustgraph/config/service/config.py b/trustgraph-flow/trustgraph/config/service/config.py index 701d7f58..ddbcc4b0 100644 --- a/trustgraph-flow/trustgraph/config/service/config.py +++ b/trustgraph-flow/trustgraph/config/service/config.py @@ -224,11 +224,7 @@ class Configuration: return ConfigResponse( version = await self.get_version(), - value = None, - directory = None, - values = None, config = config, - error = None, ) async def handle(self, msg): diff --git a/trustgraph-flow/trustgraph/config/service/service.py b/trustgraph-flow/trustgraph/config/service/service.py index 65112cd5..b84660a1 100644 --- a/trustgraph-flow/trustgraph/config/service/service.py +++ b/trustgraph-flow/trustgraph/config/service/service.py @@ -211,7 +211,6 @@ class Processor(AsyncProcessor): type = "config-error", message = str(e), ), - text=None, ) await self.config_response_producer.send( diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/collection_management.py b/trustgraph-flow/trustgraph/gateway/dispatch/collection_management.py index 9773ad4c..2fa3759d 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/collection_management.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/collection_management.py @@ -5,7 +5,7 @@ from ... messaging import TranslatorRegistry from . requestor import ServiceRequestor class CollectionManagementRequestor(ServiceRequestor): - def __init__(self, pulsar_client, consumer, subscriber, timeout=120, + def __init__(self, backend, consumer, subscriber, timeout=120, request_queue=None, response_queue=None): if request_queue is None: @@ -14,7 +14,7 @@ class CollectionManagementRequestor(ServiceRequestor): response_queue = collection_response_queue super(CollectionManagementRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, consumer_name = consumer, subscription = subscriber, request_queue=request_queue, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/config.py b/trustgraph-flow/trustgraph/gateway/dispatch/config.py index 10a0aab9..9d40e8cc 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/config.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/config.py @@ -7,7 +7,7 @@ from ... messaging import TranslatorRegistry from . requestor import ServiceRequestor class ConfigRequestor(ServiceRequestor): - def __init__(self, pulsar_client, consumer, subscriber, timeout=120, + def __init__(self, backend, consumer, subscriber, timeout=120, request_queue=None, response_queue=None): if request_queue is None: @@ -16,7 +16,7 @@ class ConfigRequestor(ServiceRequestor): response_queue = config_response_queue super(ConfigRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, consumer_name = consumer, subscription = subscriber, request_queue=request_queue, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/core_import.py b/trustgraph-flow/trustgraph/gateway/dispatch/core_import.py index b32fb7f7..af22a5b0 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/core_import.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/core_import.py @@ -11,8 +11,8 @@ logger = logging.getLogger(__name__) class CoreImport: - def __init__(self, pulsar_client): - self.pulsar_client = pulsar_client + def __init__(self, backend): + self.backend = backend async def process(self, data, error, ok, request): @@ -20,7 +20,7 @@ class CoreImport: user = request.query["user"] kr = KnowledgeRequestor( - pulsar_client = self.pulsar_client, + backend = self.backend, consumer = "api-gateway-core-import-" + str(uuid.uuid4()), subscriber = "api-gateway-core-import-" + str(uuid.uuid4()), ) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/document_load.py b/trustgraph-flow/trustgraph/gateway/dispatch/document_load.py index 7e38877c..eb68b0b1 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/document_load.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/document_load.py @@ -11,10 +11,10 @@ from . sender import ServiceSender logger = logging.getLogger(__name__) class DocumentLoad(ServiceSender): - def __init__(self, pulsar_client, queue): + def __init__(self, backend, queue): super(DocumentLoad, self).__init__( - pulsar_client = pulsar_client, + backend = backend, queue = queue, schema = Document, ) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/flow.py b/trustgraph-flow/trustgraph/gateway/dispatch/flow.py index cb641656..be91995d 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/flow.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/flow.py @@ -7,7 +7,7 @@ from ... messaging import TranslatorRegistry from . requestor import ServiceRequestor class FlowRequestor(ServiceRequestor): - def __init__(self, pulsar_client, consumer, subscriber, timeout=120, + def __init__(self, backend, consumer, subscriber, timeout=120, request_queue=None, response_queue=None): if request_queue is None: @@ -16,7 +16,7 @@ class FlowRequestor(ServiceRequestor): response_queue = flow_response_queue super(FlowRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, consumer_name = consumer, subscription = subscriber, request_queue=request_queue, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/knowledge.py b/trustgraph-flow/trustgraph/gateway/dispatch/knowledge.py index b42db648..83aefbd0 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/knowledge.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/knowledge.py @@ -10,7 +10,7 @@ from ... messaging import TranslatorRegistry from . requestor import ServiceRequestor class KnowledgeRequestor(ServiceRequestor): - def __init__(self, pulsar_client, consumer, subscriber, timeout=120, + def __init__(self, backend, consumer, subscriber, timeout=120, request_queue=None, response_queue=None): if request_queue is None: @@ -19,7 +19,7 @@ class KnowledgeRequestor(ServiceRequestor): response_queue = knowledge_response_queue super(KnowledgeRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, consumer_name = consumer, subscription = subscriber, request_queue=request_queue, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/librarian.py b/trustgraph-flow/trustgraph/gateway/dispatch/librarian.py index 8fc62d54..bbf7190e 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/librarian.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/librarian.py @@ -9,7 +9,7 @@ from ... messaging import TranslatorRegistry from . requestor import ServiceRequestor class LibrarianRequestor(ServiceRequestor): - def __init__(self, pulsar_client, consumer, subscriber, timeout=120, + def __init__(self, backend, consumer, subscriber, timeout=120, request_queue=None, response_queue=None): if request_queue is None: @@ -18,7 +18,7 @@ class LibrarianRequestor(ServiceRequestor): response_queue = librarian_response_queue super(LibrarianRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, consumer_name = consumer, subscription = subscriber, request_queue=request_queue, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/requestor.py b/trustgraph-flow/trustgraph/gateway/dispatch/requestor.py index 1acac5e5..e8f0a63e 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/requestor.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/requestor.py @@ -13,7 +13,7 @@ class ServiceRequestor: def __init__( self, - pulsar_client, + backend, request_queue, request_schema, response_queue, response_schema, subscription="api-gateway", consumer_name="api-gateway", @@ -21,12 +21,12 @@ class ServiceRequestor: ): self.pub = Publisher( - pulsar_client, request_queue, + backend, request_queue, schema=request_schema, ) self.sub = Subscriber( - pulsar_client, response_queue, + backend, response_queue, subscription, consumer_name, response_schema ) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/sender.py b/trustgraph-flow/trustgraph/gateway/dispatch/sender.py index 2435cdc1..17324b19 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/sender.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/sender.py @@ -14,12 +14,12 @@ class ServiceSender: def __init__( self, - pulsar_client, + backend, queue, schema, ): self.pub = Publisher( - pulsar_client, queue, + backend, queue, schema=schema, ) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/streamer.py b/trustgraph-flow/trustgraph/gateway/dispatch/streamer.py index 54674906..9c6d4251 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/streamer.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/streamer.py @@ -13,7 +13,7 @@ class ServiceRequestor: def __init__( self, - pulsar_client, + backend, queue, schema, handler, subscription="api-gateway", consumer_name="api-gateway", @@ -21,7 +21,7 @@ class ServiceRequestor: ): self.sub = Subscriber( - pulsar_client, queue, + backend, queue, subscription, consumer_name, schema ) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/text_load.py b/trustgraph-flow/trustgraph/gateway/dispatch/text_load.py index 36922c89..b2562938 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/text_load.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/text_load.py @@ -11,10 +11,10 @@ from . sender import ServiceSender logger = logging.getLogger(__name__) class TextLoad(ServiceSender): - def __init__(self, pulsar_client, queue): + def __init__(self, backend, queue): super(TextLoad, self).__init__( - pulsar_client = pulsar_client, + backend = backend, queue = queue, schema = TextDocument, )