mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-24 20:51:02 +02:00
Fixing incomplete porting
This commit is contained in:
parent
57c01665c5
commit
3f98e1e0ba
1 changed files with 8 additions and 8 deletions
|
|
@ -27,18 +27,18 @@ class WebSocketResponder:
|
||||||
|
|
||||||
class MessageDispatcher:
|
class MessageDispatcher:
|
||||||
|
|
||||||
def __init__(self, max_workers: int = 10, config_receiver=None, pulsar_client=None):
|
def __init__(self, max_workers: int = 10, config_receiver=None, backend=None):
|
||||||
self.max_workers = max_workers
|
self.max_workers = max_workers
|
||||||
self.semaphore = asyncio.Semaphore(max_workers)
|
self.semaphore = asyncio.Semaphore(max_workers)
|
||||||
self.active_tasks = set()
|
self.active_tasks = set()
|
||||||
self.pulsar_client = pulsar_client
|
self.backend = backend
|
||||||
|
|
||||||
# Use DispatcherManager for flow and service management
|
# Use DispatcherManager for flow and service management
|
||||||
if pulsar_client and config_receiver:
|
if backend and config_receiver:
|
||||||
self.dispatcher_manager = DispatcherManager(pulsar_client, config_receiver, prefix="rev-gateway")
|
self.dispatcher_manager = DispatcherManager(backend, config_receiver, prefix="rev-gateway")
|
||||||
else:
|
else:
|
||||||
self.dispatcher_manager = None
|
self.dispatcher_manager = None
|
||||||
logger.warning("No pulsar_client or config_receiver provided - using fallback mode")
|
logger.warning("No backend or config_receiver provided - using fallback mode")
|
||||||
|
|
||||||
# Service name mapping from websocket protocol to translator registry
|
# Service name mapping from websocket protocol to translator registry
|
||||||
self.service_mapping = {
|
self.service_mapping = {
|
||||||
|
|
@ -78,7 +78,7 @@ class MessageDispatcher:
|
||||||
|
|
||||||
try:
|
try:
|
||||||
if not self.dispatcher_manager:
|
if not self.dispatcher_manager:
|
||||||
raise RuntimeError("DispatcherManager not available - pulsar_client and config_receiver required")
|
raise RuntimeError("DispatcherManager not available - backend and config_receiver required")
|
||||||
|
|
||||||
# Use DispatcherManager for flow-based processing
|
# Use DispatcherManager for flow-based processing
|
||||||
responder = WebSocketResponder()
|
responder = WebSocketResponder()
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue