mirror of
https://github.com/trustgraph-ai/trustgraph.git
synced 2026-07-23 04:01:02 +02:00
Added config receiver
This commit is contained in:
parent
b6f52949b3
commit
556def3dc0
1 changed files with 35 additions and 1 deletions
|
|
@ -2,17 +2,22 @@ import asyncio
|
||||||
import logging
|
import logging
|
||||||
import json
|
import json
|
||||||
import sys
|
import sys
|
||||||
|
import os
|
||||||
from aiohttp import ClientSession, WSMsgType, ClientWebSocketResponse
|
from aiohttp import ClientSession, WSMsgType, ClientWebSocketResponse
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
|
import pulsar
|
||||||
|
|
||||||
from .dispatcher import MessageDispatcher
|
from .dispatcher import MessageDispatcher
|
||||||
|
from ..gateway.config.receiver import ConfigReceiver
|
||||||
|
|
||||||
logger = logging.getLogger("rev_gateway")
|
logger = logging.getLogger("rev_gateway")
|
||||||
logger.setLevel(logging.INFO)
|
logger.setLevel(logging.INFO)
|
||||||
|
|
||||||
class ReverseGateway:
|
class ReverseGateway:
|
||||||
|
|
||||||
def __init__(self, host: str = "api.trustgraph.ai", max_workers: int = 10):
|
def __init__(self, host: str = "api.trustgraph.ai", max_workers: int = 10,
|
||||||
|
pulsar_host: str = None, pulsar_api_key: str = None,
|
||||||
|
pulsar_listener: str = None):
|
||||||
self.host = host
|
self.host = host
|
||||||
self.url = f"wss://{host}/ws"
|
self.url = f"wss://{host}/ws"
|
||||||
self.max_workers = max_workers
|
self.max_workers = max_workers
|
||||||
|
|
@ -22,6 +27,27 @@ class ReverseGateway:
|
||||||
self.running = False
|
self.running = False
|
||||||
self.reconnect_delay = 3.0
|
self.reconnect_delay = 3.0
|
||||||
|
|
||||||
|
# Pulsar configuration
|
||||||
|
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
|
||||||
|
)
|
||||||
|
|
||||||
|
# Initialize config receiver
|
||||||
|
self.config_receiver = ConfigReceiver(self.pulsar_client)
|
||||||
|
|
||||||
async def connect(self) -> bool:
|
async def connect(self) -> bool:
|
||||||
try:
|
try:
|
||||||
if self.session is None:
|
if self.session is None:
|
||||||
|
|
@ -85,6 +111,10 @@ class ReverseGateway:
|
||||||
self.running = True
|
self.running = True
|
||||||
logger.info("Starting reverse gateway")
|
logger.info("Starting reverse gateway")
|
||||||
|
|
||||||
|
# Start config receiver
|
||||||
|
logger.info("Starting config receiver")
|
||||||
|
await self.config_receiver.start()
|
||||||
|
|
||||||
while self.running:
|
while self.running:
|
||||||
try:
|
try:
|
||||||
if await self.connect():
|
if await self.connect():
|
||||||
|
|
@ -113,5 +143,9 @@ class ReverseGateway:
|
||||||
await self.dispatcher.shutdown()
|
await self.dispatcher.shutdown()
|
||||||
await self.disconnect()
|
await self.disconnect()
|
||||||
|
|
||||||
|
# Close Pulsar client
|
||||||
|
if hasattr(self, 'pulsar_client'):
|
||||||
|
self.pulsar_client.close()
|
||||||
|
|
||||||
def stop(self):
|
def stop(self):
|
||||||
self.running = False
|
self.running = False
|
||||||
Loading…
Add table
Add a link
Reference in a new issue