MediaRouter / app /mcp /server.py
basyx's picture
Upload 629 files
1fed801 verified
Raw
History Blame Contribute Delete
6.48 kB
from __future__ import annotations
import argparse
import asyncio
import sys
from typing import Any, Literal, cast
import uvicorn
from mcp.server.fastmcp import FastMCP
from starlette.applications import Starlette
from starlette.routing import Mount
from app.container import Container, build_container
from app.core.config import get_settings
from app.core.logger import configure_logging
from app.mcp.prompts import register_prompts
from app.mcp.registry import MCPRegistry
from app.mcp.resources import register_resources
from app.mcp.tools.ai import register_ai_tools
from app.mcp.tools.analytics import register_analytics_tools
from app.mcp.tools.audio import register_audio_tools
from app.mcp.tools.brand import register_brand_tools
from app.mcp.tools.image import register_image_tools
from app.mcp.tools.probe import register_probe_tools
from app.mcp.tools.social import register_social_tools
from app.mcp.tools.system import register_system_tools
from app.mcp.tools.templates import register_template_tools
from app.mcp.tools.video import register_video_tools
from app.mcp.tools.whisper import register_whisper_tools
from app.mcp.tools.ytdlp import register_ytdlp_tools
from app.security.context import auth_context
from app.security.middleware import APIKeyAuthenticationMiddleware
from app.workers.cleanup_worker import CleanupWorker
MCPLogLevel = Literal["DEBUG", "INFO", "WARNING", "ERROR", "CRITICAL"]
def create_mcp_server(container: Container) -> FastMCP[Any]:
"""Create a fully registered MCP server over an existing service container."""
settings = container.settings
server: FastMCP[Any] = FastMCP(
name=settings.app_name,
instructions=(
"Use the registered tools for safe media processing. All media inputs accept URL, Base64, "
"n8n binary objects, or managed temp_path values. Read media://operations for capabilities."
),
host=settings.host,
port=settings.port,
streamable_http_path="/",
json_response=True,
stateless_http=True,
log_level=cast(MCPLogLevel, settings.log_level),
)
registry = MCPRegistry(container)
register_video_tools(server, registry)
register_audio_tools(server, registry)
register_image_tools(server, registry)
register_whisper_tools(server, registry)
register_ytdlp_tools(server, registry)
register_probe_tools(server, registry)
register_system_tools(server, registry)
register_template_tools(server, registry)
register_social_tools(server, registry)
register_brand_tools(server, registry)
register_ai_tools(server, registry)
register_analytics_tools(server, registry)
register_resources(server, registry)
register_prompts(server)
return server
async def run_server(transport: Literal["stdio", "streamable-http"] = "stdio") -> None:
"""Run standalone MCP with the same database, keys, scopes, and limits as REST."""
configure_logging(stream=sys.stderr if transport == "stdio" else None)
settings = get_settings()
container = build_container(settings)
worker = CleanupWorker(container.cleanup, settings.cleanup_interval_seconds)
server = create_mcp_server(container)
await container.security_database.initialize()
if not await container.security_database.schema_ready():
missing = ", ".join(await container.security_database.missing_schema_objects())
raise RuntimeError(
"Security schema is unavailable; apply app/security/migrations/. " f"Missing: {missing}"
)
await container.security_database.verify_execution_boundary(
expected_role=settings.security_database_role,
enforce_rls=settings.security_enforce_rls,
)
await container.api_keys.ensure_bootstrap_admin()
await container.tenants.ensure_all_api_key_principals()
await container.social.initialize()
await container.analytics.initialize(container.social.ready)
await container.social.adopt_legacy_workspaces(await container.tenants.list_principals())
await worker.start()
try:
if transport == "stdio":
context_token = None
if settings.auth_enabled:
configured_key = (
settings.mcp_stdio_api_key.get_secret_value()
if settings.mcp_stdio_api_key is not None
else ""
)
if not configured_key:
raise RuntimeError("MCP_STDIO_API_KEY is required when AUTH_ENABLED=true")
context = await container.api_keys.authenticate(configured_key)
if not context.allows("mcp:read"):
raise RuntimeError("MCP_STDIO_API_KEY is missing the mcp:read scope")
await container.api_keys.mark_used(context)
context_token = auth_context.set(context)
try:
await server.run_stdio_async()
finally:
if context_token is not None:
auth_context.reset(context_token)
else:
mcp_application = server.streamable_http_app()
application = Starlette(routes=[Mount("/mcp", app=mcp_application)])
application.state.container = container
application.state.mcp_server = server
application.add_middleware(
APIKeyAuthenticationMiddleware,
settings=settings,
api_keys=container.api_keys,
rate_limiter=container.rate_limiter,
audit=container.audit,
)
config = uvicorn.Config(
application,
host=settings.host,
port=settings.port,
log_level=settings.log_level.lower(),
)
async with server.session_manager.run():
await uvicorn.Server(config).serve()
finally:
await worker.stop()
await container.social.close()
await container.security_database.close()
def main() -> None:
"""CLI entry point for local MCP clients."""
parser = argparse.ArgumentParser(description="Enterprise Media API MCP server")
parser.add_argument(
"--transport",
choices=("stdio", "streamable-http"),
default="stdio",
help="MCP transport to run (default: stdio)",
)
arguments = parser.parse_args()
asyncio.run(run_server(arguments.transport))
if __name__ == "__main__":
main()