mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-24 20:51:02 +02:00
Fixing API mismatch
This commit is contained in:
parent
288b8b61d6
commit
caf10ff0c1
13 changed files with 24 additions and 29 deletions
|
|
@ -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):
|
||||
|
|
|
|||
|
|
@ -211,7 +211,6 @@ class Processor(AsyncProcessor):
|
|||
type = "config-error",
|
||||
message = str(e),
|
||||
),
|
||||
text=None,
|
||||
)
|
||||
|
||||
await self.config_response_producer.send(
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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()),
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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
|
||||
)
|
||||
|
|
|
|||
|
|
@ -14,12 +14,12 @@ class ServiceSender:
|
|||
|
||||
def __init__(
|
||||
self,
|
||||
pulsar_client,
|
||||
backend,
|
||||
queue, schema,
|
||||
):
|
||||
|
||||
self.pub = Publisher(
|
||||
pulsar_client, queue,
|
||||
backend, queue,
|
||||
schema=schema,
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
)
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue