diff --git a/trustgraph-flow/setup.py b/trustgraph-flow/setup.py index 1cda836d..562c5389 100644 --- a/trustgraph-flow/setup.py +++ b/trustgraph-flow/setup.py @@ -71,6 +71,7 @@ setuptools.setup( scripts=[ "scripts/agent-manager-react", "scripts/api-gateway", + "scripts/rev-gateway", "scripts/chunker-recursive", "scripts/chunker-token", "scripts/config-svc", diff --git a/trustgraph-flow/trustgraph/rev_gateway/__init__.py b/trustgraph-flow/trustgraph/rev_gateway/__init__.py index e69de29b..1be89162 100644 --- a/trustgraph-flow/trustgraph/rev_gateway/__init__.py +++ b/trustgraph-flow/trustgraph/rev_gateway/__init__.py @@ -0,0 +1 @@ +from . service import run diff --git a/trustgraph-flow/trustgraph/rev_gateway/__main__.py b/trustgraph-flow/trustgraph/rev_gateway/__main__.py index 49f77d77..70262bc8 100644 --- a/trustgraph-flow/trustgraph/rev_gateway/__main__.py +++ b/trustgraph-flow/trustgraph/rev_gateway/__main__.py @@ -1,77 +1,11 @@ -import asyncio -import argparse import logging -import sys -import os -from .service import ReverseGateway +from .service import run logging.basicConfig( level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s' ) -def parse_args(): - parser = argparse.ArgumentParser( - prog="reverse-gateway", - description="TrustGraph Reverse Gateway - WebSocket to Pulsar bridge" - ) - - parser.add_argument( - '--websocket-uri', - default=None, - help='WebSocket URI to connect to (default: wss://api.trustgraph.ai/ws or WEBSOCKET_URI env var)' - ) - - parser.add_argument( - '--max-workers', - type=int, - default=10, - help='Maximum concurrent message handlers (default: 10)' - ) - - parser.add_argument( - '--pulsar-host', - default=None, - help='Pulsar host URL (default: pulsar://pulsar:6650 or PULSAR_HOST env var)' - ) - - parser.add_argument( - '--pulsar-api-key', - default=None, - help='Pulsar API key for authentication (default: PULSAR_API_KEY env var)' - ) - - parser.add_argument( - '--pulsar-listener', - default=None, - help='Pulsar listener name' - ) - - return parser.parse_args() - -async def main(): - args = parse_args() - - gateway = ReverseGateway( - websocket_uri=args.websocket_uri, - max_workers=args.max_workers, - pulsar_host=args.pulsar_host, - pulsar_api_key=args.pulsar_api_key, - pulsar_listener=args.pulsar_listener - ) - - print(f"Starting reverse gateway:") - print(f" WebSocket URI: {gateway.url}") - print(f" Max workers: {args.max_workers}") - print(f" Pulsar host: {gateway.pulsar_host}") - - try: - await gateway.run() - except KeyboardInterrupt: - print("\nShutdown requested by user") - except Exception as e: - print(f"Fatal error: {e}") - sys.exit(1) - if __name__ == "__main__": - asyncio.run(main()) \ No newline at end of file + run() + diff --git a/trustgraph-flow/trustgraph/rev_gateway/service.py b/trustgraph-flow/trustgraph/rev_gateway/service.py index 067a8e45..e6ebda9b 100644 --- a/trustgraph-flow/trustgraph/rev_gateway/service.py +++ b/trustgraph-flow/trustgraph/rev_gateway/service.py @@ -1,4 +1,5 @@ import asyncio +import argparse import logging import json import sys @@ -173,4 +174,67 @@ class ReverseGateway: self.pulsar_client.close() def stop(self): - self.running = False \ No newline at end of file + self.running = False + +def parse_args(): + parser = argparse.ArgumentParser( + prog="reverse-gateway", + description="TrustGraph Reverse Gateway - WebSocket to Pulsar bridge" + ) + + parser.add_argument( + '--websocket-uri', + default=None, + help='WebSocket URI to connect to (default: wss://api.trustgraph.ai/ws or WEBSOCKET_URI env var)' + ) + + parser.add_argument( + '--max-workers', + type=int, + default=10, + help='Maximum concurrent message handlers (default: 10)' + ) + + parser.add_argument( + '--pulsar-host', + default=None, + help='Pulsar host URL (default: pulsar://pulsar:6650 or PULSAR_HOST env var)' + ) + + parser.add_argument( + '--pulsar-api-key', + default=None, + help='Pulsar API key for authentication (default: PULSAR_API_KEY env var)' + ) + + parser.add_argument( + '--pulsar-listener', + default=None, + help='Pulsar listener name' + ) + + return parser.parse_args() + +def run(): + args = parse_args() + + gateway = ReverseGateway( + websocket_uri=args.websocket_uri, + max_workers=args.max_workers, + pulsar_host=args.pulsar_host, + pulsar_api_key=args.pulsar_api_key, + pulsar_listener=args.pulsar_listener + ) + + print(f"Starting reverse gateway:") + print(f" WebSocket URI: {gateway.url}") + print(f" Max workers: {args.max_workers}") + print(f" Pulsar host: {gateway.pulsar_host}") + + try: + asyncio.run(gateway.run()) + except KeyboardInterrupt: + print("\nShutdown requested by user") + except Exception as e: + print(f"Fatal error: {e}") + sys.exit(1)