Fixing incomplete porting

This commit is contained in:
Cyber MacGeddon 2025-12-17 15:58:20 +00:00
parent 001e4be120
commit 57c01665c5
3 changed files with 22 additions and 27 deletions

View file

@ -234,11 +234,11 @@ class KnowledgeManager:
logger.debug(f"Graph embeddings queue: {ge_q}") logger.debug(f"Graph embeddings queue: {ge_q}")
t_pub = Publisher( t_pub = Publisher(
self.flow_config.pulsar_client, t_q, self.flow_config.pubsub, t_q,
schema=Triples, schema=Triples,
) )
ge_pub = Publisher( ge_pub = Publisher(
self.flow_config.pulsar_client, ge_q, self.flow_config.pubsub, ge_q,
schema=GraphEmbeddings schema=GraphEmbeddings
) )

View file

@ -133,12 +133,12 @@ class DispatcherManager:
async def process_core_import(self, data, error, ok, request): 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) return await ci.process(data, error, ok, request)
async def process_core_export(self, 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) return await ce.process(data, error, ok, request)
async def process_global_service(self, data, responder, params): async def process_global_service(self, data, responder, params):

View file

@ -7,10 +7,10 @@ import os
from aiohttp import ClientSession, WSMsgType, ClientWebSocketResponse from aiohttp import ClientSession, WSMsgType, ClientWebSocketResponse
from typing import Optional from typing import Optional
from urllib.parse import urlparse, urlunparse from urllib.parse import urlparse, urlunparse
import pulsar
from .dispatcher import MessageDispatcher from .dispatcher import MessageDispatcher
from ..gateway.config.receiver import ConfigReceiver from ..gateway.config.receiver import ConfigReceiver
from ..base import get_pubsub
logger = logging.getLogger("rev_gateway") logger = logging.getLogger("rev_gateway")
logger.setLevel(logging.INFO) logger.setLevel(logging.INFO)
@ -57,24 +57,19 @@ class ReverseGateway:
self.pulsar_api_key = pulsar_api_key or os.getenv("PULSAR_API_KEY", None) self.pulsar_api_key = pulsar_api_key or os.getenv("PULSAR_API_KEY", None)
self.pulsar_listener = pulsar_listener self.pulsar_listener = pulsar_listener
# Initialize Pulsar client # Create backend using factory
if self.pulsar_api_key: backend_params = {
self.pulsar_client = pulsar.Client( 'pulsar_host': self.pulsar_host,
self.pulsar_host, 'pulsar_api_key': self.pulsar_api_key,
listener_name=self.pulsar_listener, 'pulsar_listener': self.pulsar_listener,
authentication=pulsar.AuthenticationToken(self.pulsar_api_key) }
) self.backend = get_pubsub(**backend_params)
else:
self.pulsar_client = pulsar.Client(
self.pulsar_host,
listener_name=self.pulsar_listener
)
# Initialize config receiver # Initialize config receiver
self.config_receiver = ConfigReceiver(self.pulsar_client) self.config_receiver = ConfigReceiver(self.backend)
# Initialize dispatcher with config_receiver and pulsar_client - must be created after config_receiver # Initialize dispatcher with config_receiver and backend - must be created after config_receiver
self.dispatcher = MessageDispatcher(max_workers, self.config_receiver, self.pulsar_client) self.dispatcher = MessageDispatcher(max_workers, self.config_receiver, self.backend)
async def connect(self) -> bool: async def connect(self) -> bool:
try: try:
@ -171,9 +166,9 @@ class ReverseGateway:
await self.dispatcher.shutdown() await self.dispatcher.shutdown()
await self.disconnect() await self.disconnect()
# Close Pulsar client # Close backend
if hasattr(self, 'pulsar_client'): if hasattr(self, 'backend'):
self.pulsar_client.close() self.backend.close()
def stop(self): def stop(self):
self.running = False self.running = False