bootstrap.py 2.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105
  1. """Fail-closed production dependency wiring for DataOps MCP processes."""
  2. from __future__ import annotations
  3. import os
  4. from sqlalchemy import create_engine
  5. from app.core.mcp.canary import CanaryVerifier
  6. from app.core.mcp.context import ContextService
  7. from app.core.mcp.gateway import SchedulingGateway
  8. from app.core.mcp.persistence import (
  9. PostgresAuditSink,
  10. PostgresContextRepository,
  11. PostgresSchedulingPlanStore,
  12. )
  13. from app.core.orchestration.engines.kestra import KestraAdapter
  14. from app.runner.auth import TaskTokenIssuer
  15. def _required(name):
  16. value = str(os.getenv(name) or "").strip()
  17. if not value:
  18. raise ValueError(f"{name} is required")
  19. return value
  20. def _bounded_integer(name, *, default, minimum, maximum):
  21. source = str(os.getenv(name, default)).strip()
  22. try:
  23. value = int(source)
  24. except ValueError as exc:
  25. raise ValueError(f"{name} must be an integer") from exc
  26. if value < minimum or value > maximum:
  27. raise ValueError(f"{name} must be between {minimum} and {maximum}")
  28. return value
  29. def _database_engine():
  30. database_url = _required("DATAOPS_MCP_DATABASE_URL")
  31. if not database_url.startswith(("postgresql://", "postgresql+psycopg2://")):
  32. raise ValueError("DATAOPS_MCP_DATABASE_URL must use PostgreSQL")
  33. return create_engine(
  34. database_url,
  35. pool_pre_ping=True,
  36. pool_size=5,
  37. max_overflow=5,
  38. pool_recycle=300,
  39. )
  40. def _kestra_adapter():
  41. return KestraAdapter(
  42. base_url=_required("KESTRA_BASE_URL"),
  43. username=_required("KESTRA_USERNAME"),
  44. password=_required("KESTRA_PASSWORD"),
  45. tenant_id=str(os.getenv("KESTRA_TENANT_ID", "main")).strip() or "main",
  46. timeout=_bounded_integer(
  47. "KESTRA_HTTP_TIMEOUT_SECONDS",
  48. default=30,
  49. minimum=1,
  50. maximum=300,
  51. ),
  52. )
  53. def build_scheduling_gateway():
  54. engine = _database_engine()
  55. plans = PostgresSchedulingPlanStore(engine)
  56. return SchedulingGateway(
  57. plans=plans,
  58. audit=PostgresAuditSink(engine),
  59. engine=_kestra_adapter(),
  60. token_issuer=TaskTokenIssuer(
  61. _required("RUNNER_TASK_TOKEN_SECRET"),
  62. ttl_seconds=_bounded_integer(
  63. "RUNNER_TASK_TOKEN_TTL_SECONDS",
  64. default=60,
  65. minimum=1,
  66. maximum=300,
  67. ),
  68. ),
  69. )
  70. def build_context_service():
  71. engine = _database_engine()
  72. repository = PostgresContextRepository(
  73. engine,
  74. max_concurrency=_bounded_integer(
  75. "DATAOPS_MCP_CONTEXT_MAX_CONCURRENCY",
  76. default=10,
  77. minimum=1,
  78. maximum=1000,
  79. ),
  80. )
  81. return ContextService(repository)
  82. def build_canary_verifier():
  83. engine = _database_engine()
  84. return CanaryVerifier(
  85. plans=PostgresSchedulingPlanStore(engine),
  86. engine=_kestra_adapter(),
  87. )