diff --git a/.env.example b/.env.example index 7aece30..f74f9a1 100644 --- a/.env.example +++ b/.env.example @@ -63,10 +63,15 @@ COMPOSE_REDIS_PASSWORD_FILE=./.secrets/redis_password COMPOSE_GRAFANA_ADMIN_PASSWORD_FILE=./.secrets/grafana_admin_password # Connection pool configuration -DORIS_MAX_CONNECTIONS=20 +# 🔧 IMPORTANT: total physical DB connections ≈ WORKERS × DORIS_MAX_CONNECTIONS +# (before per-token pools). Keep this comfortably below the Doris user +# connection limit (max_user_connection_num, often 100). Lower this if you see +# "Reach limit of connections" errors or pool-exhaustion symptoms. +DORIS_MAX_CONNECTIONS=8 DORIS_CONNECTION_TIMEOUT=30 DORIS_HEALTH_CHECK_INTERVAL=60 -DORIS_MAX_CONNECTION_AGE=3600 +# Shorter age turns connections over faster so stale sockets are replaced sooner. +DORIS_MAX_CONNECTION_AGE=1800 # Arrow Flight SQL Configuration (Required for ADBC tools) # FE_ARROW_FLIGHT_SQL_PORT= diff --git a/.gitignore b/.gitignore index c35b99a..c7494af 100644 --- a/.gitignore +++ b/.gitignore @@ -26,3 +26,6 @@ coverage.xml tokens.json.lock /.secrets/ /tmp/ +.claude/ +.sisyphus/ +.pytest_cache/ diff --git a/doris_mcp_server/main.py b/doris_mcp_server/main.py index c996d6a..4292751 100644 --- a/doris_mcp_server/main.py +++ b/doris_mcp_server/main.py @@ -86,6 +86,7 @@ def _multiworker_environment( "DORIS_BE_WEBSERVER_PORT": str(config.database.be_webserver_port), "SERVER_HOST": host, "SERVER_PORT": str(port), + "ROUTE_PREFIX": config.route_prefix, "MCP_ALLOWED_HOSTS": ",".join(config.mcp_allowed_hosts), "MCP_ALLOWED_ORIGINS": ",".join(config.mcp_allowed_origins), "ENABLE_LEGACY_HTTP_ADAPTER": str( @@ -398,6 +399,25 @@ async def readiness_check(request: Request) -> Response: ) return JSONResponse(payload, status_code=status_code) + # Root info endpoint + async def root_info(request: Request) -> Response: + route_prefix = self.config.route_prefix + return JSONResponse( + { + "service": self.config.server_name, + "version": __version__, + "mode": "single-worker", + "mcp_initialized": True, + "route_prefix": route_prefix or None, + "endpoints": { + "health": f"{route_prefix}/health" + if route_prefix + else "/health", + "mcp": f"{route_prefix}/mcp" if route_prefix else "/mcp", + }, + } + ) + # OAuth endpoints from .auth.oauth_handlers import OAuthHandlers @@ -463,6 +483,7 @@ async def lifespan(app: Starlette) -> AsyncIterator[None]: effective_auth = get_effective_auth_config(self.config) routes = [ + Route("/", root_info, methods=["GET"]), Route("/health", health_check, methods=["GET"]), Route("/live", live_check, methods=["GET"]), Route("/ready", readiness_check, methods=["GET"]), @@ -532,6 +553,18 @@ async def mcp_app( # Handle HTTP requests if scope["type"] == "http": path = scope.get("path", "") + + # Strip the configured route prefix (reverse-proxy deployments) + # before internal routing. We deliberately do NOT pass + # uvicorn's own root_path option: it *prepends* the prefix + # to scope["path"] rather than letting us strip it here. + route_prefix = self.config.route_prefix + if route_prefix and path.startswith(route_prefix): + path = path[len(route_prefix) :] or "/" + scope = dict(scope) + scope["path"] = path + scope["root_path"] = route_prefix + self.logger.info(f"Received request for path: {path}") try: @@ -547,7 +580,8 @@ async def mcp_app( return if ( - path.rstrip("/") in {"/health", "/live", "/ready"} + path == "/" + or path.rstrip("/") in {"/health", "/live", "/ready"} or path.startswith("/auth/") or path.startswith("/token/") or path.startswith("/.well-known/") @@ -715,6 +749,13 @@ def create_arg_parser() -> argparse.ArgumentParser: help="Number of worker processes for HTTP mode (default: 1, use 0 for auto-detect CPU cores)", ) + parser.add_argument( + "--route-prefix", + type=str, + default=os.getenv("ROUTE_PREFIX", _default_config.route_prefix), + help=f"Route prefix for reverse proxy, e.g. /doris-mcp (default: {_default_config.route_prefix or 'none'})", + ) + parser.add_argument( "--doris-host", "--db-host", @@ -814,6 +855,13 @@ def cli_has(*options: str) -> bool: config.workers = args.workers _mark_source(config, "workers", "cli") + # Route prefix for reverse proxy; normalize to "/" (or empty) + # to match the normalization in DorisConfig.from_env(). + if hasattr(args, "route_prefix") and cli_has("--route-prefix"): + route_prefix = args.route_prefix.strip().strip("/") + config.route_prefix = ("/" + route_prefix) if route_prefix else "" + _mark_source(config, "route_prefix", "cli") + async def main() -> int: """Main function""" @@ -841,6 +889,8 @@ async def main() -> int: logger.info("Starting Doris MCP Server...") logger.info(f"Transport: {config.transport}") logger.info(f"Log Level: {config.logging.level}") + if config.route_prefix: + logger.info(f"Route Prefix: {config.route_prefix}") try: effective_auth = normalize_effective_auth_config( diff --git a/doris_mcp_server/multiworker_app.py b/doris_mcp_server/multiworker_app.py index 420c54f..64693c2 100644 --- a/doris_mcp_server/multiworker_app.py +++ b/doris_mcp_server/multiworker_app.py @@ -416,6 +416,8 @@ async def doris_oauth_api_refresh(request: Request) -> Response: async def root_info(request: Request) -> JSONResponse: """Root endpoint""" + route_prefix = os.getenv("ROUTE_PREFIX", "").strip().strip("/") + route_prefix = ("/" + route_prefix) if route_prefix else "" return JSONResponse( { "service": "doris-mcp-server", @@ -423,11 +425,14 @@ async def root_info(request: Request) -> JSONResponse: "worker_pid": os.getpid(), "mcp_initialized": _worker_initialized, "version": __version__, + "route_prefix": route_prefix or None, "endpoints": { - "health": "/health", - "live": "/live", - "ready": "/ready", - "mcp": MODERN_MCP_PATH, + "health": f"{route_prefix}/health" if route_prefix else "/health", + "live": f"{route_prefix}/live" if route_prefix else "/live", + "ready": f"{route_prefix}/ready" if route_prefix else "/ready", + "mcp": f"{route_prefix}{MODERN_MCP_PATH}" + if route_prefix + else MODERN_MCP_PATH, }, } ) @@ -621,6 +626,21 @@ async def app(scope: Scope, receive: Receive, send: Send) -> None: """Main ASGI app that routes requests""" path = scope.get("path", "/") + # Strip the route prefix (reverse-proxy deployments) before internal + # routing. The master process propagates ROUTE_PREFIX via + # _multiworker_environment(); read the env directly so stripping works + # even before the worker rebuilds its config. We deliberately do NOT use + # uvicorn's root_path option: it *prepends* the prefix to scope["path"] + # rather than letting us strip it here. + route_prefix = os.getenv("ROUTE_PREFIX", "").strip().strip("/") + if route_prefix: + route_prefix = "/" + route_prefix + if route_prefix and path.startswith(route_prefix): + path = path[len(route_prefix) :] or "/" + scope = dict(scope) + scope["path"] = path + scope["root_path"] = route_prefix + legacy_adapter_enabled = bool( _worker_http_transport and _worker_http_transport.legacy_adapter_enabled diff --git a/doris_mcp_server/utils/adbc_query_tools.py b/doris_mcp_server/utils/adbc_query_tools.py index 9870d21..3290388 100644 --- a/doris_mcp_server/utils/adbc_query_tools.py +++ b/doris_mcp_server/utils/adbc_query_tools.py @@ -21,10 +21,11 @@ """ import asyncio +import decimal import os import socket import time -from datetime import datetime +from datetime import date, datetime from typing import Any from ..result_limits import ( @@ -43,12 +44,20 @@ def _convert_numpy_types(obj: Any) -> Any: - """Convert numpy types to native Python types for JSON serialization""" + """Convert numpy/pandas/decimal types to native Python types for JSON serialization""" try: # Import numpy only when needed import numpy as np import pandas as pd + # 🔧 FIX: ADBC/Arrow exposes DECIMAL columns as Python decimal.Decimal. + # json.dumps() cannot serialize Decimal, so coerce to float (or str if + # too big to preserve precision). + if isinstance(obj, decimal.Decimal): + try: + return float(obj) + except (ValueError, OverflowError): + return str(obj) if isinstance(obj, np.integer): return int(obj) elif isinstance(obj, np.floating): @@ -59,12 +68,35 @@ def _convert_numpy_types(obj: Any) -> Any: return obj.tolist() elif isinstance(obj, pd.Timestamp | pd.NaT.__class__): return str(obj) - elif pd.isna(obj): - return None + elif isinstance(obj, datetime | date): + return obj.isoformat() + elif isinstance(obj, bytes | bytearray): + try: + return obj.decode("utf-8") + except UnicodeDecodeError: + return obj.hex() else: + # pd.isna() raises on some container types; only call it for scalars. + try: + if pd.isna(obj): + return None + except (TypeError, ValueError): + pass return obj except ImportError: - # If numpy/pandas not available, return as-is + # If numpy/pandas not available, still handle Decimal/datetime/bytes. + if isinstance(obj, decimal.Decimal): + try: + return float(obj) + except (ValueError, OverflowError): + return str(obj) + if isinstance(obj, datetime | date): + return obj.isoformat() + if isinstance(obj, bytes | bytearray): + try: + return obj.decode("utf-8") + except UnicodeDecodeError: + return obj.hex() return obj diff --git a/doris_mcp_server/utils/analysis_tools.py b/doris_mcp_server/utils/analysis_tools.py index 69f4932..e4ee8cd 100644 --- a/doris_mcp_server/utils/analysis_tools.py +++ b/doris_mcp_server/utils/analysis_tools.py @@ -62,78 +62,87 @@ async def get_table_summary( auth_context = get_auth_context() connection = await self.connection_manager.get_connection("query") + try: + # Get table basic information using parameterized query + table_info_sql = """ + SELECT + table_name, + table_comment, + table_rows, + create_time, + engine + FROM information_schema.tables + WHERE table_schema = DATABASE() + AND table_name = %s + """ - # Get table basic information using parameterized query - table_info_sql = """ - SELECT - table_name, - table_comment, - table_rows, - create_time, - engine - FROM information_schema.tables - WHERE table_schema = DATABASE() - AND table_name = %s - """ + table_info_result = await connection.execute( + table_info_sql, params=(table_name,), auth_context=auth_context + ) + if not table_info_result.data: + raise ValueError(f"Table {table_name} does not exist") - table_info_result = await connection.execute( - table_info_sql, params=(table_name,), auth_context=auth_context - ) - if not table_info_result.data: - raise ValueError(f"Table {table_name} does not exist") - - table_info = table_info_result.data[0] - - # Get column information using parameterized query - columns_sql = """ - SELECT - column_name, - data_type, - is_nullable, - column_comment - FROM information_schema.columns - WHERE table_schema = DATABASE() - AND table_name = %s - ORDER BY ordinal_position - """ + table_info = table_info_result.data[0] - columns_result = await connection.execute( - columns_sql, params=(table_name,), auth_context=auth_context - ) - - summary = { - "table_name": table_info["table_name"], - "comment": table_info.get("table_comment"), - "row_count": table_info.get("table_rows", 0), - "create_time": str(table_info.get("create_time")), - "engine": table_info.get("engine"), - "column_count": len(columns_result.data), - "columns": columns_result.data, - } + # Get column information using parameterized query + columns_sql = """ + SELECT + column_name, + data_type, + is_nullable, + column_comment + FROM information_schema.columns + WHERE table_schema = DATABASE() + AND table_name = %s + ORDER BY ordinal_position + """ - # Get sample data using quoted identifier - safe_sample_size = validate_integer( - sample_size, - "sample size", - minimum=0, - maximum=10000, - ) - if include_sample and safe_sample_size > 0: - quoted_table = quote_identifier(table_name, "table name") - # SQL sink audit: caller integer -> bounded validation + quoted - # identifier -> DorisConnection.execute. - sample_sql = f"SELECT * FROM {quoted_table} LIMIT {safe_sample_size}" # nosec B608 - sample_result = await connection.execute( - sample_sql, auth_context=auth_context + columns_result = await connection.execute( + columns_sql, params=(table_name,), auth_context=auth_context ) - summary["sample_data"] = sample_result.data - return summary + summary = { + "table_name": table_info["table_name"], + "comment": table_info.get("table_comment"), + "row_count": table_info.get("table_rows", 0), + "create_time": str(table_info.get("create_time")), + "engine": table_info.get("engine"), + "column_count": len(columns_result.data), + "columns": columns_result.data, + } + + # Get sample data using quoted identifier + safe_sample_size = validate_integer( + sample_size, + "sample size", + minimum=0, + maximum=10000, + ) + if include_sample and safe_sample_size > 0: + quoted_table = quote_identifier(table_name, "table name") + # SQL sink audit: caller integer -> bounded validation + quoted + # identifier -> DorisConnection.execute. + sample_sql = f"SELECT * FROM {quoted_table} LIMIT {safe_sample_size}" # nosec B608 + sample_result = await connection.execute( + sample_sql, auth_context=auth_context + ) + summary["sample_data"] = sample_result.data + + return summary + finally: + # 🔧 FIX: release the connection on every exit path (including + # exceptions) to prevent pool exhaustion. + release_connection = getattr( + self.connection_manager, "release_connection", None + ) + if connection is not None and callable(release_connection): + await release_connection("query", connection) async def analyze_column( self, table_name: str, column_name: str, analysis_type: str = "basic" ) -> dict[str, Any]: """Analyze column statistics""" + connection = None try: safe_table = quote_identifier(table_name, "table name") safe_column = quote_identifier(column_name, "column name") @@ -223,54 +232,70 @@ async def analyze_column( "column_name": column_name, "table_name": table_name, } + finally: + # 🔧 FIX: release the connection on every exit path to prevent + # pool exhaustion. + release_connection = getattr( + self.connection_manager, "release_connection", None + ) + if connection is not None and callable(release_connection): + await release_connection("query", connection) async def analyze_table_relationships( self, table_name: str, depth: int = 2 ) -> dict[str, Any]: """Analyze table relationships""" connection = await self.connection_manager.get_connection("system") + try: + # Get table basic information + table_info_sql = """ + SELECT + table_name, + table_comment, + table_rows + FROM information_schema.tables + WHERE table_schema = DATABASE() + AND table_name = %s + """ - # Get table basic information - table_info_sql = """ - SELECT - table_name, - table_comment, - table_rows - FROM information_schema.tables - WHERE table_schema = DATABASE() - AND table_name = %s - """ + auth_context = get_auth_context() + table_result = await connection.execute( + table_info_sql, + params=(table_name,), + auth_context=auth_context, + ) + if not table_result.data: + raise ValueError(f"Table {table_name} does not exist") - auth_context = get_auth_context() - table_result = await connection.execute( - table_info_sql, - params=(table_name,), - auth_context=auth_context, - ) - if not table_result.data: - raise ValueError(f"Table {table_name} does not exist") - - # Get all tables list (for analyzing potential relationships) - all_tables_sql = """ - SELECT - table_name, - table_comment - FROM information_schema.tables - WHERE table_schema = DATABASE() - AND table_type = 'BASE TABLE' - AND table_name != %s - """ + # Get all tables list (for analyzing potential relationships) + all_tables_sql = """ + SELECT + table_name, + table_comment + FROM information_schema.tables + WHERE table_schema = DATABASE() + AND table_type = 'BASE TABLE' + AND table_name != %s + """ - all_tables_result = await connection.execute( - all_tables_sql, params=(table_name,), auth_context=auth_context - ) + all_tables_result = await connection.execute( + all_tables_sql, params=(table_name,), auth_context=auth_context + ) - return { - "center_table": table_result.data[0], - "related_tables": all_tables_result.data, - "depth": depth, - "note": "Table relationship analysis based on column name similarity and business logic inference", - } + return { + "center_table": table_result.data[0], + "related_tables": all_tables_result.data, + "depth": depth, + "note": "Table relationship analysis based on column name similarity and business logic inference", + } + finally: + # 🔧 FIX: release the connection on every exit path to prevent + # pool exhaustion. + release_connection = getattr( + self.connection_manager, "release_connection", None + ) + if connection is not None and callable(release_connection): + await release_connection("system", connection) class PerformanceMonitor: @@ -284,84 +309,97 @@ async def get_performance_stats( ) -> dict[str, Any]: """Get performance statistics""" connection = await self.connection_manager.get_connection("system") + try: + if metric_type == "queries": + # Query performance metrics + stats = { + "metric_type": "queries", + "time_range": time_range, + "timestamp": datetime.now().isoformat(), + "total_queries": 0, + "avg_execution_time": 0.0, + "slow_queries": 0, + "error_queries": 0, + "note": "Query performance statistics (simulated data)", + } - if metric_type == "queries": - # Query performance metrics - stats = { - "metric_type": "queries", - "time_range": time_range, - "timestamp": datetime.now().isoformat(), - "total_queries": 0, - "avg_execution_time": 0.0, - "slow_queries": 0, - "error_queries": 0, - "note": "Query performance statistics (simulated data)", - } - - elif metric_type == "connections": - # Connection statistics - connection_metrics = await self.connection_manager.get_metrics() - stats = { - "metric_type": "connections", - "time_range": time_range, - "timestamp": datetime.now().isoformat(), - "total_connections": connection_metrics.total_connections, - "active_connections": connection_metrics.active_connections, - "idle_connections": connection_metrics.idle_connections, - "failed_connections": connection_metrics.failed_connections, - "connection_errors": connection_metrics.connection_errors, - "avg_connection_time": connection_metrics.avg_connection_time, - "last_health_check": connection_metrics.last_health_check.isoformat() - if connection_metrics.last_health_check - else None, - } - - elif metric_type == "tables": - # Table-level statistics - tables_sql = """ - SELECT - table_name, - table_rows, - data_length, - index_length, - create_time, - update_time - FROM information_schema.tables - WHERE table_schema = DATABASE() - AND table_type = 'BASE TABLE' - ORDER BY table_rows DESC - LIMIT 20 - """ + elif metric_type == "connections": + # Connection statistics + connection_metrics = await self.connection_manager.get_metrics() + stats = { + "metric_type": "connections", + "time_range": time_range, + "timestamp": datetime.now().isoformat(), + "total_connections": connection_metrics.total_connections, + "active_connections": connection_metrics.active_connections, + "idle_connections": connection_metrics.idle_connections, + "failed_connections": connection_metrics.failed_connections, + "connection_errors": connection_metrics.connection_errors, + "avg_connection_time": connection_metrics.avg_connection_time, + "last_health_check": ( + connection_metrics.last_health_check.isoformat() + if connection_metrics.last_health_check + else None + ), + } - auth_context = get_auth_context() - tables_result = await connection.execute( - tables_sql, auth_context=auth_context - ) - stats = { - "metric_type": "tables", - "time_range": time_range, - "timestamp": datetime.now().isoformat(), - "table_count": len(tables_result.data), - "tables": tables_result.data, - } + elif metric_type == "tables": + # Table-level statistics + tables_sql = """ + SELECT + table_name, + table_rows, + data_length, + index_length, + create_time, + update_time + FROM information_schema.tables + WHERE table_schema = DATABASE() + AND table_type = 'BASE TABLE' + ORDER BY table_rows DESC + LIMIT 20 + """ - elif metric_type == "system": - # System-level metrics (simulated) - stats = { - "metric_type": "system", - "time_range": time_range, - "timestamp": datetime.now().isoformat(), - "cpu_usage": 45.2, - "memory_usage": 68.5, - "disk_usage": 72.1, - "network_io": {"bytes_sent": 1024000, "bytes_received": 2048000}, - "note": "System metrics (simulated data)", - } + auth_context = get_auth_context() + tables_result = await connection.execute( + tables_sql, auth_context=auth_context + ) + stats = { + "metric_type": "tables", + "time_range": time_range, + "timestamp": datetime.now().isoformat(), + "table_count": len(tables_result.data), + "tables": tables_result.data, + } - else: - raise ValueError(f"Unsupported metric type: {metric_type}") + elif metric_type == "system": + # System-level metrics (simulated) + stats = { + "metric_type": "system", + "time_range": time_range, + "timestamp": datetime.now().isoformat(), + "cpu_usage": 45.2, + "memory_usage": 68.5, + "disk_usage": 72.1, + "network_io": { + "bytes_sent": 1024000, + "bytes_received": 2048000, + }, + "note": "System metrics (simulated data)", + } - return stats + else: + raise ValueError(f"Unsupported metric type: {metric_type}") + + return stats + finally: + # 🔧 FIX: release the connection on every exit path to prevent + # pool exhaustion. + release_connection = getattr( + self.connection_manager, "release_connection", None + ) + if connection is not None and callable(release_connection): + await release_connection("system", connection) async def get_query_history( self, limit: int = 50, order_by: str = "time" diff --git a/doris_mcp_server/utils/config.py b/doris_mcp_server/utils/config.py index 7646365..fd1163d 100644 --- a/doris_mcp_server/utils/config.py +++ b/doris_mcp_server/utils/config.py @@ -968,6 +968,10 @@ class DorisConfig: transport: str = "stdio" workers: int = 1 + # Route prefix for reverse-proxy deployments (e.g. "/doris-mcp" for an + # nginx location of /doris-mcp/). Normalized to "/" or "". + route_prefix: str = "" + # Temporary files configuration temp_files_dir: str = "tmp" # Temporary files directory for Explain and Profile outputs @@ -1606,6 +1610,13 @@ def from_env(cls, env_file: str | None = None) -> "DorisConfig": server_port = os.getenv("SERVER_PORT", "").strip() if server_port and server_port.isdigit(): config.server_port = int(server_port) + # Normalize the route prefix to "/" (or empty). Must match the + # normalization in main.update_configuration so from_env() is + # self-consistent when worker processes rebuild the config. + route_prefix = os.getenv("ROUTE_PREFIX", config.route_prefix).strip().strip("/") + config.route_prefix = ("/" + route_prefix) if route_prefix else "" + if config.route_prefix: + _mark_source(config, "route_prefix", "env") if "MCP_ALLOWED_HOSTS" in os.environ: config.mcp_allowed_hosts = [ value.strip() @@ -3196,10 +3207,19 @@ def setup_logging(self) -> None: not sys.stdout.isatty() # Not a terminal (likely piped/redirected) ) + # Include the listening port in log filenames so multi-instance + # deployments (multiple HTTP servers on different ports) do not collide. + # In stdio mode there is no listener, so keep the plain base name. + if self.config.transport == "http" and getattr(self.config, "server_port", None): + log_base_name = f"doris_mcp_server_{self.config.server_port}" + else: + log_base_name = "doris_mcp_server" + # Setup enhanced logging with cleanup functionality setup_logging( level=self.config.logging.level, log_dir=log_dir, + base_name=log_base_name, enable_console=not is_stdio_mode, # Disable console logging in stdio mode enable_file=True, enable_audit=self.config.logging.enable_audit, diff --git a/doris_mcp_server/utils/db.py b/doris_mcp_server/utils/db.py index e63d4be..32c3a64 100644 --- a/doris_mcp_server/utils/db.py +++ b/doris_mcp_server/utils/db.py @@ -895,6 +895,34 @@ async def _force_close_raw_connection( except Exception as e: self.logger.debug(f"Error force closing connection after {reason}: {e}") + async def _discard_pool_connection( + self, pool: Pool | None, raw_connection: Any, reason: str + ) -> None: + """Close a checked-out pool connection AND release its slot. + + aiomysql's Pool tracks every acquired connection in an internal + "in use" set and only removes it from that set inside release(). + Calling ensure_closed()/close() on a checked-out connection WITHOUT + release() leaves the slot permanently counted as in-use, which + exhausts the pool after ``maxsize`` such leaks. This helper closes + the underlying socket first (so the bad connection is never handed + back out) and then releases the slot so the pool can reclaim it. + """ + if not raw_connection: + return + await self._force_close_raw_connection(raw_connection, reason) + if pool is not None: + release = getattr(pool, "release", None) + if release: + try: + release(raw_connection) + except Exception as release_error: + self.logger.debug( + "pool.release() failed during discard (%s): %s", + reason, + release_error, + ) + async def _close_auth_connection(self, conn: Any) -> None: if not conn: return @@ -2182,14 +2210,12 @@ async def _warmup_pool(self) -> None: except Exception as e: self.logger.error(f"Pool warmup failed: {e}") - # Clean up any remaining connections + # Clean up any remaining connections. 🔧 FIX: release them back to + # the pool (close-then-release) instead of only ensure_closed(), + # otherwise every warmup connection held during a failure leaks a + # pool slot. for conn in warmup_connections: - try: - await conn.ensure_closed() - except Exception as cleanup_error: - self.logger.warning( - f"Failed to close warmup connection: {cleanup_error}" - ) + await self._discard_pool_connection(pool, conn, "pool warmup failure") async def _pool_health_monitor(self) -> None: """Background task to monitor pool health""" @@ -2273,27 +2299,21 @@ async def _cleanup_stale_connections(self) -> None: # Connection is healthy, release it pool.release(conn) + conn = None except TimeoutError: self.logger.debug(f"Stale connection test {i + 1} timed out") - if conn is not None: - try: - await conn.ensure_closed() - except Exception as cleanup_error: - self.logger.debug( - "Failed to close timed-out connection " - f"{i + 1}: {cleanup_error}" - ) except Exception as e: self.logger.debug(f"Stale connection test {i + 1} failed: {e}") + finally: + # 🔧 CRITICAL FIX: always close-then-release bad + # connections. ensure_closed() alone keeps the + # connection in the pool's internal "in use" set + # forever, leaking the pool slot. if conn is not None: - try: - await conn.ensure_closed() - except Exception as cleanup_error: - self.logger.debug( - "Failed to close unhealthy connection " - f"{i + 1}: {cleanup_error}" - ) + await self._discard_pool_connection( + pool, conn, "stale connection cleanup" + ) self.logger.debug( f"Stale connection cleanup completed, tested {test_count} connections" @@ -2747,6 +2767,16 @@ async def release_connection( # Check connection state before release if connection.connection.closed: self.logger.debug(f"Connection already closed for session {session_id}") + # Still notify the pool so its internal "in use" counters stay + # correct: a checked-out closed connection that is never + # released() permanently occupies a pool slot. release() + # safely drops closed connections instead of re-queuing them. + try: + self.pool.release(connection.connection) + except Exception as release_error: + self.logger.debug( + f"Pool release of closed connection failed: {release_error}" + ) return # 🔧 FIX: Simplified release operation without thread wrapper diff --git a/doris_mcp_server/utils/dependency_analysis_tools.py b/doris_mcp_server/utils/dependency_analysis_tools.py index a7b5940..2ffe7d0 100644 --- a/doris_mcp_server/utils/dependency_analysis_tools.py +++ b/doris_mcp_server/utils/dependency_analysis_tools.py @@ -68,6 +68,7 @@ async def analyze_data_flow_dependencies( Returns: Comprehensive dependency analysis results """ + connection = None try: start_time = time.time() connection = await self.connection_manager.get_connection("query") @@ -125,6 +126,14 @@ async def analyze_data_flow_dependencies( "error": str(e), "analysis_timestamp": datetime.now().isoformat() } + finally: + # 🔧 FIX: release the connection on every exit path (success, + # error-dict return, and exception) to prevent pool exhaustion. + release_connection = getattr( + self.connection_manager, "release_connection", None + ) + if connection is not None and callable(release_connection): + await release_connection("query", connection) # ==================== Private Helper Methods ==================== diff --git a/doris_mcp_server/utils/logger.py b/doris_mcp_server/utils/logger.py index 1d257e1..78a1e8c 100644 --- a/doris_mcp_server/utils/logger.py +++ b/doris_mcp_server/utils/logger.py @@ -29,6 +29,7 @@ import os import sys import threading +import time from datetime import datetime, timedelta from pathlib import Path from typing import Any, Literal @@ -48,13 +49,20 @@ def __init__( style: Literal["%", "{", "$"] = "%", ) -> None: if fmt is None: - fmt = "%(asctime)s.%(msecs)03d %(level_aligned)s %(name)s:%(lineno)d - %(message)s" + fmt = "%(asctime)s.%(msecs)03d %(level_aligned)s [user=%(user)s token=%(token_id)s ip=%(client_ip)s] %(name)s:%(lineno)d - %(message)s" if datefmt is None: datefmt = "%Y-%m-%d %H:%M:%S" super().__init__(fmt, datefmt, style) def format(self, record: logging.LogRecord) -> str: """Format log record with enhanced information and proper alignment""" + # Ensure request-auth fields always resolve. RequestAuthContextFilter + # normally sets them on every record; setdefault is a safety net for + # records that somehow bypassed the filter, so format() never KeyErrors. + record.__dict__.setdefault("user", "-") + record.__dict__.setdefault("token_id", "-") + record.__dict__.setdefault("client_ip", "-") + # Add process info if available if hasattr(record, 'process') and record.process: record.process_info = f"[PID:{record.process}]" @@ -77,6 +85,103 @@ def format(self, record: logging.LogRecord) -> str: return super().format(record) +class RequestAuthContextFilter(logging.Filter): + """Inject the current request's user identity into every log record. + + Reads ``AuthContext`` from the ``mcp_auth_context_var`` ContextVar (set + per-request by the auth middleware / HTTP entrypoints) and exposes + ``user`` / ``token_id`` / ``client_ip`` on the record so the formatter + can include them. Falls back to '-' when no request is in flight + (stdio mode, startup logs, background threads). + + Attach to a Handler (not a Logger) so it runs for records emitted by any + submodule. ``Filter.filter()`` runs before ``Formatter.format()``, so the + ``%(user)s`` / ``%(token_id)s`` / ``%(client_ip)s`` placeholders resolve. + """ + + def filter(self, record: logging.LogRecord) -> bool: + try: + # Lazy import avoids a circular dependency (security imports + # logger helpers for its own logging). + from .security import mcp_auth_context_var + + ctx = mcp_auth_context_var.get() + except Exception: + ctx = None + record.user = getattr(ctx, "user_id", "") or "-" + record.token_id = getattr(ctx, "token_id", "") or "-" + record.client_ip = getattr(ctx, "client_ip", "") or "-" + return True + + +class TimestampRotatingFileHandler(logging.handlers.RotatingFileHandler): + """Size-based rotating handler whose backups are stamped with the rollover time. + + Rotates by file size exactly like ``RotatingFileHandler``, but instead of + producing ``file.log.1``, ``file.log.2`` ... it renames the current file to + ``file_.log`` on each rollover and keeps only the most + recent ``backupCount`` backups (oldest-by-mtime are pruned). + """ + + _TS_FMT = "%Y%m%d_%H%M%S" + + def _rotated_path(self, ts: str) -> str: + base = self.baseFilename + if base.endswith(".log"): + return f"{base[:-len('.log')]}_{ts}.log" + return f"{base}_{ts}" + + def _backup_files(self) -> list[str]: + """List rotated backup files owned by this handler (excludes the live file).""" + base = self.baseFilename + parent = os.path.dirname(base) or "." + if base.endswith(".log"): + stem = os.path.basename(base[:-len(".log")]) + else: + stem = os.path.basename(base) + prefix = f"{stem}_" + files: list[str] = [] + try: + for name in os.listdir(parent): + if name.startswith(prefix) and name.endswith(".log"): + files.append(os.path.join(parent, name)) + except OSError: + return [] + return files + + def doRollover(self) -> None: + """Rotate the current log to a timestamped backup, pruning old backups.""" + if self.stream: + self.stream.close() + self.stream = None + + ts = datetime.now().strftime(self._TS_FMT) + rotated = self._rotated_path(ts) + # Disambiguate same-second rollovers with a millisecond suffix + if os.path.exists(rotated): + ms = int(time.time() * 1000) % 1000 + rotated = self._rotated_path(f"{ts}_{ms:03d}") + + if os.path.exists(self.baseFilename): + try: + os.rename(self.baseFilename, rotated) + except OSError: + pass + + # Prune oldest backups beyond backupCount (keep newest by mtime) + if self.backupCount > 0: + backups = self._backup_files() + backups.sort(key=lambda p: os.path.getmtime(p)) + for stale in backups[:-self.backupCount]: + try: + os.remove(stale) + except OSError: + pass + + if not self.delay: + self.stream = self._open() + + class LevelBasedFileHandler(logging.Handler): """Custom handler that writes different log levels to different files""" @@ -109,7 +214,7 @@ def _setup_level_handlers(self) -> None: for level, filename in level_files.items(): file_path = self.log_dir / f"{self.base_name}_{filename}" - handler = logging.handlers.RotatingFileHandler( + handler = TimestampRotatingFileHandler( file_path, maxBytes=self.max_bytes, backupCount=self.backup_count, @@ -138,7 +243,8 @@ def close(self) -> None: class LogCleanupManager: """Log file cleanup manager for automatic maintenance""" - def __init__(self, log_dir: str, max_age_days: int = 30, cleanup_interval_hours: int = 24): + def __init__(self, log_dir: str, max_age_days: int = 30, cleanup_interval_hours: int = 24, + cleanup_at_startup: bool = False): """ Initialize log cleanup manager. @@ -146,10 +252,15 @@ def __init__(self, log_dir: str, max_age_days: int = 30, cleanup_interval_hours: log_dir: Directory containing log files max_age_days: Maximum age of log files in days (default: 30 days) cleanup_interval_hours: Cleanup interval in hours (default: 24 hours) + cleanup_at_startup: If True, run a cleanup pass immediately when the + scheduler starts. When False (default) the first cleanup waits one + full interval, so logs left from a previous run are not purged on + boot. """ self.log_dir = Path(log_dir) self.max_age_days = max_age_days self.cleanup_interval_hours = cleanup_interval_hours + self.cleanup_at_startup = cleanup_at_startup self.cleanup_thread: threading.Thread | None = None self.stop_event = threading.Event() self.logger: logging.Logger | None = None @@ -178,10 +289,21 @@ def stop_cleanup_scheduler(self) -> None: self.logger.info("Log cleanup scheduler stopped") def _cleanup_loop(self) -> None: - """Background loop for periodic cleanup""" + """Background loop for periodic cleanup. + + When ``cleanup_at_startup`` is False (the default), the loop waits one + full interval before its first cleanup pass, so logs from a previous + run are not purged the moment the server boots. + """ + first_pass = True while not self.stop_event.is_set(): try: - self.cleanup_old_logs() + # Skip the immediate cleanup on the very first pass unless + # cleanup_at_startup was explicitly requested. + if first_pass and not self.cleanup_at_startup: + first_pass = False + else: + self.cleanup_old_logs() # Sleep for the specified interval, but check stop event every 60 seconds for _ in range(self.cleanup_interval_hours * 60): # Convert hours to minutes if self.stop_event.wait(60): # Wait 60 seconds or until stop event @@ -300,6 +422,7 @@ def __init__(self) -> None: def setup_logging(self, level: str = "INFO", log_dir: str = "logs", + base_name: str = "doris_mcp_server", enable_console: bool = True, enable_file: bool = True, enable_audit: bool = True, @@ -315,6 +438,8 @@ def setup_logging(self, Args: level: Base logging level (DEBUG, INFO, WARNING, ERROR, CRITICAL) log_dir: Directory for log files + base_name: Base name for log files (e.g. "doris_mcp_server_3000" + when running multiple instances on different ports) enable_console: Enable console output enable_file: Enable file logging enable_audit: Enable audit logging @@ -359,8 +484,9 @@ def setup_logging(self, console_handler = logging.StreamHandler(sys.stdout) console_handler.setLevel(getattr(logging, level.upper())) console_handler.addFilter(_sensitive_data_filter) + console_handler.addFilter(RequestAuthContextFilter()) console_formatter = TimestampedFormatter( - fmt="%(asctime)s.%(msecs)03d %(level_aligned)s %(name)s - %(message)s", + fmt="%(asctime)s.%(msecs)03d %(level_aligned)s [user=%(user)s token=%(token_id)s ip=%(client_ip)s] %(name)s - %(message)s", datefmt="%Y-%m-%d %H:%M:%S" ) console_handler.setFormatter(console_formatter) @@ -370,18 +496,21 @@ def setup_logging(self, if enable_file: level_handler = LevelBasedFileHandler( log_dir=str(self.log_dir), - base_name="doris_mcp_server", + base_name=base_name, max_bytes=max_file_size, backup_count=backup_count ) level_handler.setLevel(logging.DEBUG) # Accept all levels level_handler.addFilter(_sensitive_data_filter) + # Attach to the composite handler (not the per-level sub-handlers): + # sub-handlers receive records via emit(), which bypasses their filters. + level_handler.addFilter(RequestAuthContextFilter()) handlers.append(level_handler) # Combined application log (all levels in one file) if enable_file: - app_log_file = self.log_dir / "doris_mcp_server_all.log" - app_handler = logging.handlers.RotatingFileHandler( + app_log_file = self.log_dir / f"{base_name}_all.log" + app_handler = TimestampRotatingFileHandler( app_log_file, maxBytes=max_file_size, backupCount=backup_count, @@ -389,13 +518,14 @@ def setup_logging(self, ) app_handler.setLevel(getattr(logging, level.upper())) app_handler.addFilter(_sensitive_data_filter) + app_handler.addFilter(RequestAuthContextFilter()) app_formatter = TimestampedFormatter() app_handler.setFormatter(app_formatter) handlers.append(app_handler) # Audit logger (separate from main logging) if enable_audit: - audit_file_path = audit_file or str(self.log_dir / "doris_mcp_server_audit.log") + audit_file_path = audit_file or str(self.log_dir / f"{base_name}_audit.log") audit_logger = logging.getLogger("audit") audit_logger.setLevel(logging.INFO) @@ -403,15 +533,16 @@ def setup_logging(self, for handler in audit_logger.handlers[:]: audit_logger.removeHandler(handler) - audit_handler = logging.handlers.RotatingFileHandler( + audit_handler = TimestampRotatingFileHandler( audit_file_path, maxBytes=max_file_size, backupCount=backup_count, encoding='utf-8' ) audit_handler.addFilter(_sensitive_data_filter) + audit_handler.addFilter(RequestAuthContextFilter()) audit_formatter = TimestampedFormatter( - fmt="%(asctime)s.%(msecs)03d [AUDIT] %(name)s - %(message)s", + fmt="%(asctime)s.%(msecs)03d [AUDIT] [user=%(user)s token=%(token_id)s ip=%(client_ip)s] %(name)s - %(message)s", datefmt="%Y-%m-%d %H:%M:%S" ) audit_handler.setFormatter(audit_formatter) @@ -563,6 +694,7 @@ def shutdown(self) -> None: def setup_logging(level: str = "INFO", log_dir: str = "logs", + base_name: str = "doris_mcp_server", enable_console: bool = True, enable_file: bool = True, enable_audit: bool = True, @@ -578,6 +710,7 @@ def setup_logging(level: str = "INFO", Args: level: Logging level (DEBUG, INFO, WARNING, ERROR, CRITICAL) log_dir: Directory for log files + base_name: Base name for log files (e.g. "doris_mcp_server_3000") enable_console: Enable console output enable_file: Enable file logging enable_audit: Enable audit logging @@ -591,6 +724,7 @@ def setup_logging(level: str = "INFO", _logger_manager.setup_logging( level=level, log_dir=log_dir, + base_name=base_name, enable_console=enable_console, enable_file=enable_file, enable_audit=enable_audit, diff --git a/doris_mcp_server/utils/schema_extractor.py b/doris_mcp_server/utils/schema_extractor.py index 07e4dda..88d264c 100644 --- a/doris_mcp_server/utils/schema_extractor.py +++ b/doris_mcp_server/utils/schema_extractor.py @@ -1844,7 +1844,12 @@ async def get_table_comment_async( db_name: str | None = None, catalog_name: str | None = None, ) -> str: - """Async version: get the comment for a table.""" + """Async version: get the comment for a table. + + Some older Doris versions do not expose TABLE_COMMENT in + information_schema.tables ("Unknown column 'table_comment'"). Fall back + to SHOW TABLE STATUS, which has a Comment column on all versions. + """ try: effective_db = db_name or self.db_name effective_catalog = catalog_name or self.catalog_name @@ -1868,23 +1873,79 @@ async def get_table_comment_async( AND TABLE_NAME = %s """ - result = await self._execute_query_with_catalog_async( - query, - effective_db, - effective_catalog, - params=(effective_db, table_name), - ) + try: + result = await self._execute_query_with_catalog_async( + query, + effective_db, + effective_catalog, + params=(effective_db, table_name), + ) + except Exception as query_error: + # Older Doris versions lack TABLE_COMMENT; fall back to + # SHOW TABLE STATUS. + if ( + "unknown column" in str(query_error).lower() + or "table_comment" in str(query_error).lower() + ): + logger.debug( + "information_schema.tables has no TABLE_COMMENT on this " + "Doris version; falling back to SHOW TABLE STATUS for %s", + table_name, + ) + result = await self._get_table_comment_via_show_status_async( + table_name, effective_db, effective_catalog + ) + else: + raise + if not result or not result[0]: self._raise_if_doris_oauth_table_not_visible( table_name, effective_db, effective_catalog ) return "" - return result[0].get("TABLE_COMMENT", "") or "" + row = result[0] + # SHOW TABLE STATUS returns "Comment"; information_schema returns + # "TABLE_COMMENT". + return row.get("TABLE_COMMENT") or row.get("Comment") or "" except Exception as e: - logger.error(f"Failed to get table comment asynchronously: {e}") + logger.debug(f"Failed to get table comment asynchronously: {e}") self._reraise_if_doris_oauth_metadata_error(e) return "" + async def _get_table_comment_via_show_status_async( + self, + table_name: str, + db_name: str, + catalog_name: str | None, + ) -> list: + """Fallback comment retrieval via SHOW TABLE STATUS (cross-version safe).""" + quoted_db = quote_identifier(db_name, "database name") + query = f"SHOW TABLE STATUS FROM {quoted_db} LIKE '{table_name}'" + try: + result = await self._execute_query_with_catalog_async( + query, db_name, catalog_name + ) + # Normalize the Comment column name to TABLE_COMMENT for uniform + # parsing. + if result: + normalized = [] + for row in result: + if isinstance(row, dict): + normalized.append( + { + **row, + "TABLE_COMMENT": row.get( + "Comment", row.get("TABLE_COMMENT", "") + ), + } + ) + else: + normalized.append(row) + return normalized + return result + except Exception: + return [] + async def get_column_comments_async( self, table_name: str, diff --git a/examples/examples.py b/examples/examples.py new file mode 100644 index 0000000..0031e2f --- /dev/null +++ b/examples/examples.py @@ -0,0 +1,31 @@ +import asyncio +from doris_mcp_client.client import DorisUnifiedClient, DorisClientConfig +# HTTP 模式 +async def example_http(): + config = DorisClientConfig.http("http://localhost:3000/doris-mcp/mcp", timeout=60) + client = DorisUnifiedClient(config) + async def operations(client: DorisUnifiedClient): + # 列出所有工具 + tools = await client.list_all_tools() + print(f"可用工具: {[t.name for t in tools]}") + # 获取数据库列表 + db_list = await client.get_database_list() + print(f"数据库列表: {db_list}") + # 执行 SQL 查询 + result = await client.execute_sql("SELECT COUNT(*) FROM internal.ssb.customer") + print(f"查询结果: {result}") + # 获取表结构 + schema = await client.get_table_schema("customer", "ssb") + print(f"表结构: {schema}") + await client.connect_and_run(operations) +# Stdio 模式 +async def example_stdio(): + config = DorisClientConfig.stdio( + "doris-mcp-server", ["--transport", "stdio"] + ) + client = DorisUnifiedClient(config) + async def operations(client: DorisUnifiedClient): + result = await client.execute_sql("SELECT 1 AS test") + print(f"结果: {result}") + await client.connect_and_run(operations) +asyncio.run(example_http()) \ No newline at end of file diff --git a/start_server.sh b/start_server.sh index 21da3d7..b3afb6b 100755 --- a/start_server.sh +++ b/start_server.sh @@ -46,8 +46,10 @@ find . -type d -name "__pycache__" -exec rm -rf {} + 2>/dev/null || true find . -type f -name "*.pyc" -delete 2>/dev/null || true echo -e "${CYAN}Cleaning temporary files...${NC}" rm -rf .pytest_cache 2>/dev/null || true -echo -e "${CYAN}Cleaning log files...${NC}" -find ./logs -type f -name "*.log" -delete 2>/dev/null || true + +# NOTE: Old log files are intentionally NOT deleted on startup so that logs +# from a previous run are preserved. The server's LogCleanupManager handles +# aged logs on its own schedule (see doris_mcp_server/utils/logger.py). # Create necessary directories mkdir -p logs @@ -109,6 +111,9 @@ export SERVER_HOST="${SERVER_HOST:-${MCP_HOST:-127.0.0.1}}" export MCP_HOST="${SERVER_HOST}" export SERVER_PORT="${SERVER_PORT:-3000}" # Changed from MCP_PORT to SERVER_PORT export WORKERS="${WORKERS:-1}" +# Optional reverse-proxy route prefix, e.g. ROUTE_PREFIX=doris-mcp behind an +# nginx `location /doris-mcp/`. Empty (default) serves from the root path. +export ROUTE_PREFIX="${ROUTE_PREFIX:-}" export ALLOWED_ORIGINS="${ALLOWED_ORIGINS:-*}" export LOG_LEVEL="${LOG_LEVEL:-info}" export MCP_ALLOW_CREDENTIALS="${MCP_ALLOW_CREDENTIALS:-false}" @@ -117,17 +122,28 @@ export MCP_ALLOW_CREDENTIALS="${MCP_ALLOW_CREDENTIALS:-false}" export MCP_DEBUG_ADAPTER="true" export PYTHONPATH="$(pwd):$PYTHONPATH" +# Build the URL path prefix for the echoed endpoints below. +ROUTE_PATH_PREFIX="" +ROUTE_PREFIX_ARGS=() +if [ -n "${ROUTE_PREFIX}" ]; then + ROUTE_PATH_PREFIX="/${ROUTE_PREFIX#/}" + ROUTE_PREFIX_ARGS=(--route-prefix "${ROUTE_PREFIX}") +fi + echo -e "${GREEN}Starting MCP server (Streamable HTTP mode)...${NC}" -echo -e "${YELLOW}Service will run on http://${SERVER_HOST}:${SERVER_PORT}/mcp${NC}" -echo -e "${YELLOW}Liveness: http://${SERVER_HOST}:${SERVER_PORT}/live${NC}" -echo -e "${YELLOW}Readiness: http://${SERVER_HOST}:${SERVER_PORT}/ready${NC}" -echo -e "${YELLOW}MCP Endpoint: http://${SERVER_HOST}:${SERVER_PORT}/mcp${NC}" -echo -e "${YELLOW}Local access: http://localhost:${SERVER_PORT}/mcp${NC}" +echo -e "${YELLOW}Service will run on http://${SERVER_HOST}:${SERVER_PORT}${ROUTE_PATH_PREFIX}/mcp${NC}" +echo -e "${YELLOW}Liveness: http://${SERVER_HOST}:${SERVER_PORT}${ROUTE_PATH_PREFIX}/live${NC}" +echo -e "${YELLOW}Readiness: http://${SERVER_HOST}:${SERVER_PORT}${ROUTE_PATH_PREFIX}/ready${NC}" +echo -e "${YELLOW}MCP Endpoint: http://${SERVER_HOST}:${SERVER_PORT}${ROUTE_PATH_PREFIX}/mcp${NC}" +echo -e "${YELLOW}Local access: http://localhost:${SERVER_PORT}${ROUTE_PATH_PREFIX}/mcp${NC}" echo -e "${YELLOW}Workers: ${WORKERS}${NC}" +if [ -n "${ROUTE_PREFIX}" ]; then + echo -e "${YELLOW}Route Prefix: /${ROUTE_PREFIX#/}${NC}" +fi echo -e "${YELLOW}Use Ctrl+C to stop the service${NC}" # Start the server in HTTP mode (Streamable HTTP) -python -m doris_mcp_server.main --transport http --host "${SERVER_HOST}" --port "${SERVER_PORT}" --workers "${WORKERS}" +python -m doris_mcp_server.main --transport http --host "${SERVER_HOST}" --port "${SERVER_PORT}" --workers "${WORKERS}" "${ROUTE_PREFIX_ARGS[@]}" # Check exit status if [ $? -ne 0 ]; then @@ -139,5 +155,5 @@ fi echo -e "${YELLOW}Tip: If the page displays abnormally, please clear your browser cache or use incognito mode${NC}" echo -e "${YELLOW}Chrome browser clear cache shortcut: Ctrl+Shift+Del (Windows) or Cmd+Shift+Del (Mac)${NC}" echo -e "${CYAN}For testing HTTP endpoints, you can use:${NC}" -echo -e "${CYAN} curl --fail http://127.0.0.1:${SERVER_PORT}/live${NC}" -echo -e "${CYAN} curl --fail http://127.0.0.1:${SERVER_PORT}/ready${NC}" +echo -e "${CYAN} curl --fail http://127.0.0.1:${SERVER_PORT}${ROUTE_PATH_PREFIX}/live${NC}" +echo -e "${CYAN} curl --fail http://127.0.0.1:${SERVER_PORT}${ROUTE_PATH_PREFIX}/ready${NC}"