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.rules import (
  20. PostgresRulePlanRepository,
  21. RulePlanExecutor,
  22. SqlRulePlanAdapter,
  23. )
  24. def _required(name):
  25. value = str(os.environ.get(name, "")).strip()
  26. if not value:
  27. raise ValueError(f"{name} is required")
  28. return value
  29. def _integer(name, default, minimum, maximum):
  30. try:
  31. value = int(os.environ.get(name, str(default)))
  32. except (TypeError, ValueError) as exc:
  33. raise ValueError(f"{name} must be an integer") from exc
  34. if value < minimum or value > maximum:
  35. raise ValueError(f"{name} must be between {minimum} and {maximum}")
  36. return value
  37. @dataclass(frozen=True)
  38. class RunnerSettings:
  39. runtime: DataSourceRuntimeConfig = field(repr=False)
  40. task_token_secret: str = field(repr=False)
  41. allowed_http_hosts: frozenset = field(default_factory=frozenset)
  42. task_token_ttl_seconds: int = 60
  43. max_query_rows: int = 1000
  44. def runner_settings_from_env():
  45. task_token_secret = _required("RUNNER_TASK_TOKEN_SECRET")
  46. worker_count = _integer("RUNNER_WORKERS", 2, 1, 8)
  47. runtime = DataSourceRuntimeConfig.from_mapping(
  48. {
  49. "platform_database_url": _required("DATABASE_URL"),
  50. "neo4j_uri": _required("NEO4J_URI"),
  51. "neo4j_user": _required("NEO4J_USER"),
  52. "neo4j_password": _required("NEO4J_PASSWORD"),
  53. "credential_master_key": _required(
  54. "DATASOURCE_CREDENTIAL_MASTER_KEY"
  55. ),
  56. "credential_key_version": _required(
  57. "DATASOURCE_CREDENTIAL_KEY_VERSION"
  58. ),
  59. "certificate_dir": os.environ.get(
  60. "DATASOURCE_CERT_DIR",
  61. "/etc/dataops-platform/datasource-certs",
  62. ),
  63. "pool_size": _integer("RUNNER_DATASOURCE_POOL_SIZE", 1, 1, 3),
  64. "max_overflow": _integer(
  65. "RUNNER_DATASOURCE_MAX_OVERFLOW", 1, 0, 3
  66. ),
  67. "pool_timeout": _integer(
  68. "RUNNER_DATASOURCE_POOL_TIMEOUT", 10, 1, 60
  69. ),
  70. "pool_recycle": _integer(
  71. "RUNNER_DATASOURCE_POOL_RECYCLE", 1800, 60, 86400
  72. ),
  73. "idle_ttl": _integer(
  74. "RUNNER_DATASOURCE_POOL_IDLE_TTL", 900, 60, 86400
  75. ),
  76. "max_idle_pools": _integer(
  77. "RUNNER_DATASOURCE_MAX_IDLE_POOLS", 4, 1, 20
  78. ),
  79. "query_timeout": _integer(
  80. "RUNNER_DATASOURCE_QUERY_TIMEOUT", 30, 1, 300
  81. ),
  82. "worker_count": worker_count,
  83. "connection_budget": _integer(
  84. "RUNNER_DATASOURCE_CONNECTION_BUDGET", 32, 1, 200
  85. ),
  86. }
  87. )
  88. allowed_hosts = frozenset(
  89. host.strip().lower()
  90. for host in os.environ.get("RUNNER_HTTP_ALLOWED_HOSTS", "").split(",")
  91. if host.strip()
  92. )
  93. return RunnerSettings(
  94. runtime=runtime,
  95. task_token_secret=task_token_secret,
  96. allowed_http_hosts=allowed_hosts,
  97. task_token_ttl_seconds=_integer(
  98. "RUNNER_TASK_TOKEN_TTL_SECONDS", 60, 1, 300
  99. ),
  100. max_query_rows=_integer("RUNNER_MAX_QUERY_ROWS", 1000, 1, 10000),
  101. )
  102. def build_runner_application(settings=None):
  103. settings = settings or runner_settings_from_env()
  104. runtime = build_standalone_data_source_runtime(settings.runtime)
  105. query_executor = SqlQueryExecutor(
  106. runtime.manager,
  107. max_rows=settings.max_query_rows,
  108. )
  109. write_executor = SqlExecuteExecutor(runtime.manager)
  110. sql_rule_adapter = SqlRulePlanAdapter(
  111. query_executor=query_executor,
  112. write_executor=write_executor,
  113. )
  114. rule_executor = RulePlanExecutor(
  115. PostgresRulePlanRepository(runtime.platform_engine),
  116. adapters={
  117. "sql_pushdown": sql_rule_adapter,
  118. "quality_check": sql_rule_adapter,
  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