From 57c01665c5ae3639b1492cc60d86407879020cf4 Mon Sep 17 00:00:00 2001 From: Cyber MacGeddon Date: Wed, 17 Dec 2025 15:58:20 +0000 Subject: [PATCH] Fixing incomplete porting --- trustgraph-flow/trustgraph/cores/knowledge.py | 4 +- .../trustgraph/gateway/dispatch/manager.py | 4 +- .../trustgraph/rev_gateway/service.py | 41 ++++++++----------- 3 files changed, 22 insertions(+), 27 deletions(-) diff --git a/trustgraph-flow/trustgraph/cores/knowledge.py b/trustgraph-flow/trustgraph/cores/knowledge.py index 449f1c3b..0d5c3d82 100644 --- a/trustgraph-flow/trustgraph/cores/knowledge.py +++ b/trustgraph-flow/trustgraph/cores/knowledge.py @@ -234,11 +234,11 @@ class KnowledgeManager: logger.debug(f"Graph embeddings queue: {ge_q}") t_pub = Publisher( - self.flow_config.pulsar_client, t_q, + self.flow_config.pubsub, t_q, schema=Triples, ) ge_pub = Publisher( - self.flow_config.pulsar_client, ge_q, + self.flow_config.pubsub, ge_q, schema=GraphEmbeddings ) diff --git a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py index 7978c3a0..0766e232 100644 --- a/trustgraph-flow/trustgraph/gateway/dispatch/manager.py +++ b/trustgraph-flow/trustgraph/gateway/dispatch/manager.py @@ -133,12 +133,12 @@ class DispatcherManager: async def process_core_import(self, data, error, ok, request): - ci = CoreImport(self.pulsar_client) + ci = CoreImport(self.backend) return await ci.process(data, error, ok, request) async def process_core_export(self, data, error, ok, request): - ce = CoreExport(self.pulsar_client) + ce = CoreExport(self.backend) return await ce.process(data, error, ok, request) async def process_global_service(self, data, responder, params): diff --git a/trustgraph-flow/trustgraph/rev_gateway/service.py b/trustgraph-flow/trustgraph/rev_gateway/service.py index c8e78af2..cc905172 100644 --- a/trustgraph-flow/trustgraph/rev_gateway/service.py +++ b/trustgraph-flow/trustgraph/rev_gateway/service.py @@ -7,10 +7,10 @@ import os from aiohttp import ClientSession, WSMsgType, ClientWebSocketResponse from typing import Optional from urllib.parse import urlparse, urlunparse -import pulsar from .dispatcher import MessageDispatcher from ..gateway.config.receiver import ConfigReceiver +from ..base import get_pubsub logger = logging.getLogger("rev_gateway") logger.setLevel(logging.INFO) @@ -56,25 +56,20 @@ class ReverseGateway: self.pulsar_host = pulsar_host or os.getenv("PULSAR_HOST", "pulsar://pulsar:6650") self.pulsar_api_key = pulsar_api_key or os.getenv("PULSAR_API_KEY", None) self.pulsar_listener = pulsar_listener - - # Initialize Pulsar client - if self.pulsar_api_key: - self.pulsar_client = pulsar.Client( - self.pulsar_host, - listener_name=self.pulsar_listener, - authentication=pulsar.AuthenticationToken(self.pulsar_api_key) - ) - else: - self.pulsar_client = pulsar.Client( - self.pulsar_host, - listener_name=self.pulsar_listener - ) - + + # Create backend using factory + backend_params = { + 'pulsar_host': self.pulsar_host, + 'pulsar_api_key': self.pulsar_api_key, + 'pulsar_listener': self.pulsar_listener, + } + self.backend = get_pubsub(**backend_params) + # Initialize config receiver - self.config_receiver = ConfigReceiver(self.pulsar_client) - - # Initialize dispatcher with config_receiver and pulsar_client - must be created after config_receiver - self.dispatcher = MessageDispatcher(max_workers, self.config_receiver, self.pulsar_client) + self.config_receiver = ConfigReceiver(self.backend) + + # Initialize dispatcher with config_receiver and backend - must be created after config_receiver + self.dispatcher = MessageDispatcher(max_workers, self.config_receiver, self.backend) async def connect(self) -> bool: try: @@ -170,10 +165,10 @@ class ReverseGateway: self.running = False await self.dispatcher.shutdown() await self.disconnect() - - # Close Pulsar client - if hasattr(self, 'pulsar_client'): - self.pulsar_client.close() + + # Close backend + if hasattr(self, 'backend'): + self.backend.close() def stop(self): self.running = False