From 001e4be120171109506a8ea07add32b5f3c23a7a Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 17 Dec 2025 15:54:12 +0000 Subject: [PATCH] Fixing some more pulsar API mismatches --- trustgraph-flow/trustgraph/gateway/dispatch/agent.py | 4 ++-- .../trustgraph/gateway/dispatch/document_embeddings_import.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/document_rag.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/embeddings.py | 4 ++-- .../trustgraph/gateway/dispatch/entity_contexts_import.py | 4 ++-- .../trustgraph/gateway/dispatch/graph_embeddings_import.py | 4 ++-- .../trustgraph/gateway/dispatch/graph_embeddings_query.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/graph_rag.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/mcp_tool.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/objects_import.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/prompt.py | 4 ++-- .../trustgraph/gateway/dispatch/structured_diag.py | 4 ++-- .../trustgraph/gateway/dispatch/structured_query.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/triples_import.py | 4 ++-- trustgraph-flow/trustgraph/gateway/dispatch/triples_query.py | 4 ++-- 17 files changed, 34 insertions(+), 34 deletions(-) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/agent.py b/trustgraph-flow/trustgraph/gateway/dispatch/agent.py index 1a5e8299..8867956d 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/agent.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/agent.py @@ -6,12 +6,12 @@ from . requestor import ServiceRequestor class AgentRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(AgentRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=AgentRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/document_embeddings_import.py b/trustgraph-flow/trustgraph/gateway/dispatch/document_embeddings_import.py index 7ec2f595..bd5f9666 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/document_embeddings_import.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/document_embeddings_import.py @@ -15,7 +15,7 @@ logger = logging.getLogger(__name__) class DocumentEmbeddingsImport: def __init__( - self, ws, running, pulsar_client, queue + self, ws, running, backend, queue ): self.ws = ws @@ -23,7 +23,7 @@ class DocumentEmbeddingsImport: self.translator = DocumentEmbeddingsTranslator() self.publisher = Publisher( - pulsar_client, topic = queue, schema = DocumentEmbeddings + backend, topic = queue, schema = DocumentEmbeddings ) async def start(self): diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/document_rag.py b/trustgraph-flow/trustgraph/gateway/dispatch/document_rag.py index a7f3634e..83b3cb9a 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/document_rag.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/document_rag.py @@ -6,12 +6,12 @@ from . requestor import ServiceRequestor class DocumentRagRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(DocumentRagRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=DocumentRagQuery, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/embeddings.py b/trustgraph-flow/trustgraph/gateway/dispatch/embeddings.py index 47146e57..6c1b55ba 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/embeddings.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/embeddings.py @@ -6,12 +6,12 @@ from . requestor import ServiceRequestor class EmbeddingsRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(EmbeddingsRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=EmbeddingsRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/entity_contexts_import.py b/trustgraph-flow/trustgraph/gateway/dispatch/entity_contexts_import.py index c76f1612..6e01a5ca 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/entity_contexts_import.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/entity_contexts_import.py @@ -16,14 +16,14 @@ logger = logging.getLogger(__name__) class EntityContextsImport: def __init__( - self, ws, running, pulsar_client, queue + self, ws, running, backend, queue ): self.ws = ws self.running = running self.publisher = Publisher( - pulsar_client, topic = queue, schema = EntityContexts + backend, topic = queue, schema = EntityContexts ) async def start(self): diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/graph_embeddings_import.py b/trustgraph-flow/trustgraph/gateway/dispatch/graph_embeddings_import.py index ee3d88ef..8abf5e9c 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/graph_embeddings_import.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/graph_embeddings_import.py @@ -16,14 +16,14 @@ logger = logging.getLogger(__name__) class GraphEmbeddingsImport: def __init__( - self, ws, running, pulsar_client, queue + self, ws, running, backend, queue ): self.ws = ws self.running = running self.publisher = Publisher( - pulsar_client, topic = queue, schema = GraphEmbeddings + backend, topic = queue, schema = GraphEmbeddings ) async def start(self): diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/graph_embeddings_query.py b/trustgraph-flow/trustgraph/gateway/dispatch/graph_embeddings_query.py index f5be06fb..a7bb1bd8 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/graph_embeddings_query.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/graph_embeddings_query.py @@ -6,12 +6,12 @@ from . requestor import ServiceRequestor class GraphEmbeddingsQueryRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(GraphEmbeddingsQueryRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=GraphEmbeddingsRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/graph_rag.py b/trustgraph-flow/trustgraph/gateway/dispatch/graph_rag.py index a15a1aee..a0299a43 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/graph_rag.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/graph_rag.py @@ -6,12 +6,12 @@ from . requestor import ServiceRequestor class GraphRagRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(GraphRagRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=GraphRagQuery, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/mcp_tool.py b/trustgraph-flow/trustgraph/gateway/dispatch/mcp_tool.py index da2a7bb0..a5f9398e 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/mcp_tool.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/mcp_tool.py @@ -6,12 +6,12 @@ from . requestor import ServiceRequestor class McpToolRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(McpToolRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=ToolRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py b/trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py index 3cf5684a..3a6314f2 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/nlp_query.py @@ -5,12 +5,12 @@ from . requestor import ServiceRequestor class NLPQueryRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(NLPQueryRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=QuestionToStructuredQueryRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/objects_import.py b/trustgraph-flow/trustgraph/gateway/dispatch/objects_import.py index bc0c1b85..fc982b69 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/objects_import.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/objects_import.py @@ -15,14 +15,14 @@ logger = logging.getLogger(__name__) class ObjectsImport: def __init__( - self, ws, running, pulsar_client, queue + self, ws, running, backend, queue ): self.ws = ws self.running = running self.publisher = Publisher( - pulsar_client, topic = queue, schema = ExtractedObject + backend, topic = queue, schema = ExtractedObject ) async def start(self): diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py b/trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py index 2f2535a9..fb8dc81d 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/objects_query.py @@ -5,12 +5,12 @@ from . requestor import ServiceRequestor class ObjectsQueryRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(ObjectsQueryRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=ObjectsQueryRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/prompt.py b/trustgraph-flow/trustgraph/gateway/dispatch/prompt.py index 5c316cf6..23017733 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/prompt.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/prompt.py @@ -8,12 +8,12 @@ from . requestor import ServiceRequestor class PromptRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(PromptRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=PromptRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/structured_diag.py b/trustgraph-flow/trustgraph/gateway/dispatch/structured_diag.py index 8dae646d..895b55be 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/structured_diag.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/structured_diag.py @@ -5,12 +5,12 @@ from . requestor import ServiceRequestor class StructuredDiagRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(StructuredDiagRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=StructuredDataDiagnosisRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/structured_query.py b/trustgraph-flow/trustgraph/gateway/dispatch/structured_query.py index f08ef038..9a9fbb6a 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/structured_query.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/structured_query.py @@ -5,12 +5,12 @@ from . requestor import ServiceRequestor class StructuredQueryRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(StructuredQueryRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=StructuredQueryRequest, diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/triples_import.py b/trustgraph-flow/trustgraph/gateway/dispatch/triples_import.py index 520a9cbc..6bb46975 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/triples_import.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/triples_import.py @@ -16,14 +16,14 @@ logger = logging.getLogger(__name__) class TriplesImport: def __init__( - self, ws, running, pulsar_client, queue + self, ws, running, backend, queue ): self.ws = ws self.running = running self.publisher = Publisher( - pulsar_client, topic = queue, schema = Triples + backend, topic = queue, schema = Triples ) async def start(self): diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/triples_query.py b/trustgraph-flow/trustgraph/gateway/dispatch/triples_query.py index d2def9c1..6b306139 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/triples_query.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/triples_query.py @@ -6,12 +6,12 @@ from . requestor import ServiceRequestor class TriplesQueryRequestor(ServiceRequestor): def __init__( - self, pulsar_client, request_queue, response_queue, timeout, + self, backend, request_queue, response_queue, timeout, consumer, subscriber, ): super(TriplesQueryRequestor, self).__init__( - pulsar_client=pulsar_client, + backend=backend, request_queue=request_queue, response_queue=response_queue, request_schema=TriplesQueryRequest,