mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-24 20:51:02 +02:00
Fixing some more pulsar API mismatches
This commit is contained in:
parent
8e92f01c97
commit
001e4be120
17 changed files with 34 additions and 34 deletions
|
|
@ -6,12 +6,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class AgentRequestor(ServiceRequestor):
|
class AgentRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(AgentRequestor, self).__init__(
|
super(AgentRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=AgentRequest,
|
request_schema=AgentRequest,
|
||||||
|
|
|
||||||
|
|
@ -15,7 +15,7 @@ logger = logging.getLogger(__name__)
|
||||||
class DocumentEmbeddingsImport:
|
class DocumentEmbeddingsImport:
|
||||||
|
|
||||||
def __init__(
|
def __init__(
|
||||||
self, ws, running, pulsar_client, queue
|
self, ws, running, backend, queue
|
||||||
):
|
):
|
||||||
|
|
||||||
self.ws = ws
|
self.ws = ws
|
||||||
|
|
@ -23,7 +23,7 @@ class DocumentEmbeddingsImport:
|
||||||
self.translator = DocumentEmbeddingsTranslator()
|
self.translator = DocumentEmbeddingsTranslator()
|
||||||
|
|
||||||
self.publisher = Publisher(
|
self.publisher = Publisher(
|
||||||
pulsar_client, topic = queue, schema = DocumentEmbeddings
|
backend, topic = queue, schema = DocumentEmbeddings
|
||||||
)
|
)
|
||||||
|
|
||||||
async def start(self):
|
async def start(self):
|
||||||
|
|
|
||||||
|
|
@ -6,12 +6,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class DocumentRagRequestor(ServiceRequestor):
|
class DocumentRagRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(DocumentRagRequestor, self).__init__(
|
super(DocumentRagRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=DocumentRagQuery,
|
request_schema=DocumentRagQuery,
|
||||||
|
|
|
||||||
|
|
@ -6,12 +6,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class EmbeddingsRequestor(ServiceRequestor):
|
class EmbeddingsRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(EmbeddingsRequestor, self).__init__(
|
super(EmbeddingsRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=EmbeddingsRequest,
|
request_schema=EmbeddingsRequest,
|
||||||
|
|
|
||||||
|
|
@ -16,14 +16,14 @@ logger = logging.getLogger(__name__)
|
||||||
class EntityContextsImport:
|
class EntityContextsImport:
|
||||||
|
|
||||||
def __init__(
|
def __init__(
|
||||||
self, ws, running, pulsar_client, queue
|
self, ws, running, backend, queue
|
||||||
):
|
):
|
||||||
|
|
||||||
self.ws = ws
|
self.ws = ws
|
||||||
self.running = running
|
self.running = running
|
||||||
|
|
||||||
self.publisher = Publisher(
|
self.publisher = Publisher(
|
||||||
pulsar_client, topic = queue, schema = EntityContexts
|
backend, topic = queue, schema = EntityContexts
|
||||||
)
|
)
|
||||||
|
|
||||||
async def start(self):
|
async def start(self):
|
||||||
|
|
|
||||||
|
|
@ -16,14 +16,14 @@ logger = logging.getLogger(__name__)
|
||||||
class GraphEmbeddingsImport:
|
class GraphEmbeddingsImport:
|
||||||
|
|
||||||
def __init__(
|
def __init__(
|
||||||
self, ws, running, pulsar_client, queue
|
self, ws, running, backend, queue
|
||||||
):
|
):
|
||||||
|
|
||||||
self.ws = ws
|
self.ws = ws
|
||||||
self.running = running
|
self.running = running
|
||||||
|
|
||||||
self.publisher = Publisher(
|
self.publisher = Publisher(
|
||||||
pulsar_client, topic = queue, schema = GraphEmbeddings
|
backend, topic = queue, schema = GraphEmbeddings
|
||||||
)
|
)
|
||||||
|
|
||||||
async def start(self):
|
async def start(self):
|
||||||
|
|
|
||||||
|
|
@ -6,12 +6,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class GraphEmbeddingsQueryRequestor(ServiceRequestor):
|
class GraphEmbeddingsQueryRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(GraphEmbeddingsQueryRequestor, self).__init__(
|
super(GraphEmbeddingsQueryRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=GraphEmbeddingsRequest,
|
request_schema=GraphEmbeddingsRequest,
|
||||||
|
|
|
||||||
|
|
@ -6,12 +6,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class GraphRagRequestor(ServiceRequestor):
|
class GraphRagRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(GraphRagRequestor, self).__init__(
|
super(GraphRagRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=GraphRagQuery,
|
request_schema=GraphRagQuery,
|
||||||
|
|
|
||||||
|
|
@ -6,12 +6,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class McpToolRequestor(ServiceRequestor):
|
class McpToolRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(McpToolRequestor, self).__init__(
|
super(McpToolRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=ToolRequest,
|
request_schema=ToolRequest,
|
||||||
|
|
|
||||||
|
|
@ -5,12 +5,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class NLPQueryRequestor(ServiceRequestor):
|
class NLPQueryRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(NLPQueryRequestor, self).__init__(
|
super(NLPQueryRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=QuestionToStructuredQueryRequest,
|
request_schema=QuestionToStructuredQueryRequest,
|
||||||
|
|
|
||||||
|
|
@ -15,14 +15,14 @@ logger = logging.getLogger(__name__)
|
||||||
class ObjectsImport:
|
class ObjectsImport:
|
||||||
|
|
||||||
def __init__(
|
def __init__(
|
||||||
self, ws, running, pulsar_client, queue
|
self, ws, running, backend, queue
|
||||||
):
|
):
|
||||||
|
|
||||||
self.ws = ws
|
self.ws = ws
|
||||||
self.running = running
|
self.running = running
|
||||||
|
|
||||||
self.publisher = Publisher(
|
self.publisher = Publisher(
|
||||||
pulsar_client, topic = queue, schema = ExtractedObject
|
backend, topic = queue, schema = ExtractedObject
|
||||||
)
|
)
|
||||||
|
|
||||||
async def start(self):
|
async def start(self):
|
||||||
|
|
|
||||||
|
|
@ -5,12 +5,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class ObjectsQueryRequestor(ServiceRequestor):
|
class ObjectsQueryRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(ObjectsQueryRequestor, self).__init__(
|
super(ObjectsQueryRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=ObjectsQueryRequest,
|
request_schema=ObjectsQueryRequest,
|
||||||
|
|
|
||||||
|
|
@ -8,12 +8,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class PromptRequestor(ServiceRequestor):
|
class PromptRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(PromptRequestor, self).__init__(
|
super(PromptRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=PromptRequest,
|
request_schema=PromptRequest,
|
||||||
|
|
|
||||||
|
|
@ -5,12 +5,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class StructuredDiagRequestor(ServiceRequestor):
|
class StructuredDiagRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(StructuredDiagRequestor, self).__init__(
|
super(StructuredDiagRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=StructuredDataDiagnosisRequest,
|
request_schema=StructuredDataDiagnosisRequest,
|
||||||
|
|
|
||||||
|
|
@ -5,12 +5,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class StructuredQueryRequestor(ServiceRequestor):
|
class StructuredQueryRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(StructuredQueryRequestor, self).__init__(
|
super(StructuredQueryRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=StructuredQueryRequest,
|
request_schema=StructuredQueryRequest,
|
||||||
|
|
|
||||||
|
|
@ -16,14 +16,14 @@ logger = logging.getLogger(__name__)
|
||||||
class TriplesImport:
|
class TriplesImport:
|
||||||
|
|
||||||
def __init__(
|
def __init__(
|
||||||
self, ws, running, pulsar_client, queue
|
self, ws, running, backend, queue
|
||||||
):
|
):
|
||||||
|
|
||||||
self.ws = ws
|
self.ws = ws
|
||||||
self.running = running
|
self.running = running
|
||||||
|
|
||||||
self.publisher = Publisher(
|
self.publisher = Publisher(
|
||||||
pulsar_client, topic = queue, schema = Triples
|
backend, topic = queue, schema = Triples
|
||||||
)
|
)
|
||||||
|
|
||||||
async def start(self):
|
async def start(self):
|
||||||
|
|
|
||||||
|
|
@ -6,12 +6,12 @@ from . requestor import ServiceRequestor
|
||||||
|
|
||||||
class TriplesQueryRequestor(ServiceRequestor):
|
class TriplesQueryRequestor(ServiceRequestor):
|
||||||
def __init__(
|
def __init__(
|
||||||
self, pulsar_client, request_queue, response_queue, timeout,
|
self, backend, request_queue, response_queue, timeout,
|
||||||
consumer, subscriber,
|
consumer, subscriber,
|
||||||
):
|
):
|
||||||
|
|
||||||
super(TriplesQueryRequestor, self).__init__(
|
super(TriplesQueryRequestor, self).__init__(
|
||||||
pulsar_client=pulsar_client,
|
backend=backend,
|
||||||
request_queue=request_queue,
|
request_queue=request_queue,
|
||||||
response_queue=response_queue,
|
response_queue=response_queue,
|
||||||
request_schema=TriplesQueryRequest,
|
request_schema=TriplesQueryRequest,
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue