| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105 |
- """Fail-closed production dependency wiring for DataOps MCP processes."""
- from __future__ import annotations
- import os
- from sqlalchemy import create_engine
- from app.core.mcp.canary import CanaryVerifier
- from app.core.mcp.context import ContextService
- from app.core.mcp.gateway import SchedulingGateway
- from app.core.mcp.persistence import (
- PostgresAuditSink,
- PostgresContextRepository,
- PostgresSchedulingPlanStore,
- )
- from app.core.orchestration.engines.kestra import KestraAdapter
- from app.runner.auth import TaskTokenIssuer
- def _required(name):
- value = str(os.getenv(name) or "").strip()
- if not value:
- raise ValueError(f"{name} is required")
- return value
- def _bounded_integer(name, *, default, minimum, maximum):
- source = str(os.getenv(name, default)).strip()
- try:
- value = int(source)
- except ValueError as exc:
- raise ValueError(f"{name} must be an integer") from exc
- if value < minimum or value > maximum:
- raise ValueError(f"{name} must be between {minimum} and {maximum}")
- return value
- def _database_engine():
- database_url = _required("DATAOPS_MCP_DATABASE_URL")
- if not database_url.startswith(("postgresql://", "postgresql+psycopg2://")):
- raise ValueError("DATAOPS_MCP_DATABASE_URL must use PostgreSQL")
- return create_engine(
- database_url,
- pool_pre_ping=True,
- pool_size=5,
- max_overflow=5,
- pool_recycle=300,
- )
- def _kestra_adapter():
- return KestraAdapter(
- base_url=_required("KESTRA_BASE_URL"),
- username=_required("KESTRA_USERNAME"),
- password=_required("KESTRA_PASSWORD"),
- tenant_id=str(os.getenv("KESTRA_TENANT_ID", "main")).strip() or "main",
- timeout=_bounded_integer(
- "KESTRA_HTTP_TIMEOUT_SECONDS",
- default=30,
- minimum=1,
- maximum=300,
- ),
- )
- def build_scheduling_gateway():
- engine = _database_engine()
- plans = PostgresSchedulingPlanStore(engine)
- return SchedulingGateway(
- plans=plans,
- audit=PostgresAuditSink(engine),
- engine=_kestra_adapter(),
- token_issuer=TaskTokenIssuer(
- _required("RUNNER_TASK_TOKEN_SECRET"),
- ttl_seconds=_bounded_integer(
- "RUNNER_TASK_TOKEN_TTL_SECONDS",
- default=60,
- minimum=1,
- maximum=300,
- ),
- ),
- )
- def build_context_service():
- engine = _database_engine()
- repository = PostgresContextRepository(
- engine,
- max_concurrency=_bounded_integer(
- "DATAOPS_MCP_CONTEXT_MAX_CONCURRENCY",
- default=10,
- minimum=1,
- maximum=1000,
- ),
- )
- return ContextService(repository)
- def build_canary_verifier():
- engine = _database_engine()
- return CanaryVerifier(
- plans=PostgresSchedulingPlanStore(engine),
- engine=_kestra_adapter(),
- )
|