bootstrap.py 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153
  1. """Environment parsing and dependency wiring for the standalone Runner."""
  2. from __future__ import annotations
  3. import os
  4. from dataclasses import dataclass, field
  5. from app.core.data_source.runtime import (
  6. DataSourceRuntimeConfig,
  7. build_standalone_data_source_runtime,
  8. )
  9. from app.runner.api import create_runner_app
  10. from app.runner.auth import TaskTokenVerifier
  11. from app.runner.ledger import PostgresTaskLedger
  12. from app.runner.nodes import (
  13. GovernedHttpExecutor,
  14. NodeRegistry,
  15. RestrictedPythonExecutor,
  16. SqlExecuteExecutor,
  17. SqlQueryExecutor,
  18. )
  19. from app.runner.rule_sql import (
  20. SqlGlotQualityPlanAdapter,
  21. SqlGlotRulePlanAdapter,
  22. )
  23. from app.runner.rules import (
  24. PostgresRulePlanRepository,
  25. RulePlanExecutor,
  26. )
  27. def _required(name):
  28. value = str(os.environ.get(name, "")).strip()
  29. if not value:
  30. raise ValueError(f"{name} is required")
  31. return value
  32. def _integer(name, default, minimum, maximum):
  33. try:
  34. value = int(os.environ.get(name, str(default)))
  35. except (TypeError, ValueError) as exc:
  36. raise ValueError(f"{name} must be an integer") from exc
  37. if value < minimum or value > maximum:
  38. raise ValueError(f"{name} must be between {minimum} and {maximum}")
  39. return value
  40. @dataclass(frozen=True)
  41. class RunnerSettings:
  42. runtime: DataSourceRuntimeConfig = field(repr=False)
  43. task_token_secret: str = field(repr=False)
  44. allowed_http_hosts: frozenset = field(default_factory=frozenset)
  45. task_token_ttl_seconds: int = 60
  46. max_query_rows: int = 1000
  47. def runner_settings_from_env():
  48. task_token_secret = _required("RUNNER_TASK_TOKEN_SECRET")
  49. worker_count = _integer("RUNNER_WORKERS", 2, 1, 8)
  50. runtime = DataSourceRuntimeConfig.from_mapping(
  51. {
  52. "platform_database_url": _required("DATABASE_URL"),
  53. "neo4j_uri": _required("NEO4J_URI"),
  54. "neo4j_user": _required("NEO4J_USER"),
  55. "neo4j_password": _required("NEO4J_PASSWORD"),
  56. "credential_master_key": _required(
  57. "DATASOURCE_CREDENTIAL_MASTER_KEY"
  58. ),
  59. "credential_key_version": _required(
  60. "DATASOURCE_CREDENTIAL_KEY_VERSION"
  61. ),
  62. "certificate_dir": os.environ.get(
  63. "DATASOURCE_CERT_DIR",
  64. "/etc/dataops-platform/datasource-certs",
  65. ),
  66. "pool_size": _integer("RUNNER_DATASOURCE_POOL_SIZE", 1, 1, 3),
  67. "max_overflow": _integer(
  68. "RUNNER_DATASOURCE_MAX_OVERFLOW", 1, 0, 3
  69. ),
  70. "pool_timeout": _integer(
  71. "RUNNER_DATASOURCE_POOL_TIMEOUT", 10, 1, 60
  72. ),
  73. "pool_recycle": _integer(
  74. "RUNNER_DATASOURCE_POOL_RECYCLE", 1800, 60, 86400
  75. ),
  76. "idle_ttl": _integer(
  77. "RUNNER_DATASOURCE_POOL_IDLE_TTL", 900, 60, 86400
  78. ),
  79. "max_idle_pools": _integer(
  80. "RUNNER_DATASOURCE_MAX_IDLE_POOLS", 4, 1, 20
  81. ),
  82. "query_timeout": _integer(
  83. "RUNNER_DATASOURCE_QUERY_TIMEOUT", 30, 1, 300
  84. ),
  85. "worker_count": worker_count,
  86. "connection_budget": _integer(
  87. "RUNNER_DATASOURCE_CONNECTION_BUDGET", 32, 1, 200
  88. ),
  89. }
  90. )
  91. allowed_hosts = frozenset(
  92. host.strip().lower()
  93. for host in os.environ.get("RUNNER_HTTP_ALLOWED_HOSTS", "").split(",")
  94. if host.strip()
  95. )
  96. return RunnerSettings(
  97. runtime=runtime,
  98. task_token_secret=task_token_secret,
  99. allowed_http_hosts=allowed_hosts,
  100. task_token_ttl_seconds=_integer(
  101. "RUNNER_TASK_TOKEN_TTL_SECONDS", 60, 1, 300
  102. ),
  103. max_query_rows=_integer("RUNNER_MAX_QUERY_ROWS", 1000, 1, 10000),
  104. )
  105. def build_runner_application(settings=None):
  106. settings = settings or runner_settings_from_env()
  107. runtime = build_standalone_data_source_runtime(settings.runtime)
  108. query_executor = SqlQueryExecutor(
  109. runtime.manager,
  110. max_rows=settings.max_query_rows,
  111. )
  112. write_executor = SqlExecuteExecutor(runtime.manager)
  113. sql_rule_adapter = SqlGlotRulePlanAdapter(runtime.manager)
  114. rule_executor = RulePlanExecutor(
  115. PostgresRulePlanRepository(runtime.platform_engine),
  116. adapters={
  117. "sql_pushdown": sql_rule_adapter,
  118. "quality_check": SqlGlotQualityPlanAdapter(),
  119. },
  120. )
  121. registry = NodeRegistry(
  122. {
  123. "sql.query": query_executor,
  124. "sql.execute": write_executor,
  125. "python": RestrictedPythonExecutor({}),
  126. "http": GovernedHttpExecutor(
  127. allowed_hosts=settings.allowed_http_hosts
  128. ),
  129. "rule.apply": rule_executor,
  130. "quality.check": rule_executor,
  131. }
  132. )
  133. application = create_runner_app(
  134. verifier=TaskTokenVerifier(settings.task_token_secret),
  135. ledger=PostgresTaskLedger(runtime.platform_engine),
  136. registry=registry,
  137. )
  138. application.extensions["dataops_runner_runtime"] = runtime
  139. application.extensions["dataops_runner_settings"] = settings
  140. return application