api.py 8.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210
  1. """Minimal HTTP surface for the independent DataOps Runner."""
  2. import hashlib
  3. import json
  4. from flask import Flask, jsonify, request
  5. from app.runner.auth import TaskTokenInvalid
  6. from app.runner.nodes import NodeExecutionError
  7. def create_runner_app(*, verifier, ledger, registry):
  8. app = Flask("dataops-runner")
  9. app.config["MAX_CONTENT_LENGTH"] = 2 * 1024 * 1024
  10. @app.get("/health")
  11. def health():
  12. return jsonify({"status": "ok"})
  13. def finish_safely(jti, **outcome):
  14. try:
  15. ledger.finish(jti, **outcome)
  16. return True
  17. except Exception:
  18. app.logger.error("runner task ledger unavailable")
  19. return False
  20. @app.post("/v1/tasks/execute")
  21. def execute_task():
  22. payload = request.get_json(silent=True)
  23. if not isinstance(payload, dict):
  24. return jsonify({"error": "invalid task request"}), 400
  25. node = payload.get("node")
  26. parameters = payload.get("parameters", {})
  27. if not isinstance(node, dict) or not isinstance(parameters, dict):
  28. return jsonify({"error": "invalid task request"}), 400
  29. try:
  30. claims = verifier.verify(payload.get("task_token"), node=node)
  31. except TaskTokenInvalid:
  32. return jsonify({"error": "task token is invalid"}), 401
  33. binding = {
  34. "task_uid": claims.task_uid,
  35. "dataflow_uid": claims.dataflow_uid,
  36. "deployment_id": claims.deployment_id,
  37. "environment": claims.environment,
  38. "workflow_version": claims.workflow_version,
  39. "correlation_id": claims.correlation_id,
  40. "node_id": claims.node_id,
  41. "node_type": claims.node_type,
  42. "data_source_uid": node.get("data_source_uid"),
  43. "idempotency_key": (node.get("idempotency") or {}).get("key"),
  44. }
  45. try:
  46. claimed = ledger.claim(
  47. claims.jti,
  48. binding,
  49. expires_at=claims.expires_at,
  50. )
  51. except Exception:
  52. app.logger.error("runner task ledger unavailable")
  53. return jsonify({"error": "task ledger is unavailable"}), 503
  54. if not claimed:
  55. try:
  56. existing = ledger.get(claims.jti)
  57. except Exception:
  58. return jsonify({"error": "task ledger is unavailable"}), 503
  59. if (
  60. existing is None
  61. or any(
  62. existing.binding.get(key) != value
  63. for key, value in binding.items()
  64. )
  65. ):
  66. return jsonify({"error": "task token already consumed"}), 409
  67. if existing.status == "running":
  68. recovered = None
  69. if node.get("type") in {"rule.apply", "quality.check"}:
  70. try:
  71. recovered = registry.replay_task(
  72. node,
  73. correlation_id=claims.correlation_id,
  74. deployment_id=claims.deployment_id,
  75. task_jti=claims.jti,
  76. )
  77. except NodeExecutionError:
  78. return jsonify(
  79. {"error": "task ledger is unavailable"}
  80. ), 503
  81. if isinstance(recovered, dict):
  82. replay_body = {
  83. "task_uid": claims.task_uid,
  84. "correlation_id": claims.correlation_id,
  85. "result": recovered,
  86. }
  87. output_artifact = recovered.get("output_artifact")
  88. if isinstance(output_artifact, str):
  89. replay_body["output_artifact"] = output_artifact
  90. if not finish_safely(
  91. claims.jti,
  92. status="success",
  93. commit_outcome=recovered.get(
  94. "commit_outcome",
  95. "not_applicable",
  96. ),
  97. safe_detail="task response recovered",
  98. replay_http_status=200,
  99. replay_body=replay_body,
  100. ):
  101. return jsonify(
  102. {"error": "task ledger is unavailable"}
  103. ), 503
  104. response = jsonify(replay_body)
  105. response.headers["X-Idempotent-Replay"] = "true"
  106. return response
  107. response = jsonify(
  108. {
  109. "task_uid": claims.task_uid,
  110. "correlation_id": claims.correlation_id,
  111. "status": "running",
  112. }
  113. )
  114. response.headers["Retry-After"] = "2"
  115. return response, 202
  116. if (
  117. existing.replay_body is None
  118. or existing.replay_http_status is None
  119. or existing.replay_digest is None
  120. ):
  121. return jsonify({"error": "task token already consumed"}), 409
  122. encoded = json.dumps(
  123. existing.replay_body,
  124. sort_keys=True,
  125. separators=(",", ":"),
  126. ensure_ascii=False,
  127. ).encode("utf-8")
  128. if (
  129. hashlib.sha256(encoded).hexdigest()
  130. != existing.replay_digest
  131. ):
  132. return jsonify({"error": "task ledger is unavailable"}), 503
  133. response = jsonify(existing.replay_body)
  134. response.headers["X-Idempotent-Replay"] = "true"
  135. return response, existing.replay_http_status
  136. try:
  137. result = registry.execute(
  138. node,
  139. parameters,
  140. write_authorized=claims.write_authorized,
  141. correlation_id=claims.correlation_id,
  142. dataflow_uid=claims.dataflow_uid,
  143. deployment_id=claims.deployment_id,
  144. environment=claims.environment,
  145. workflow_version=claims.workflow_version,
  146. node_id=claims.node_id,
  147. task_jti=claims.jti,
  148. )
  149. except NodeExecutionError as exc:
  150. recorded = finish_safely(
  151. claims.jti,
  152. status=(
  153. "unknown"
  154. if exc.commit_outcome == "unknown"
  155. else "failed"
  156. ),
  157. commit_outcome=exc.commit_outcome,
  158. safe_detail=str(exc),
  159. )
  160. if not recorded:
  161. return jsonify({"error": "task ledger is unavailable"}), 503
  162. return jsonify({"error": str(exc)}), 400
  163. except Exception:
  164. finish_safely(
  165. claims.jti,
  166. status="failed",
  167. safe_detail="task execution failed",
  168. )
  169. app.logger.exception("runner task failed")
  170. return jsonify({"error": "task execution failed"}), 500
  171. commit_outcome = (
  172. result.get("commit_outcome", "not_applicable")
  173. if isinstance(result, dict)
  174. else "not_applicable"
  175. )
  176. response = {
  177. "task_uid": claims.task_uid,
  178. "correlation_id": claims.correlation_id,
  179. "result": result,
  180. }
  181. if (
  182. isinstance(result, dict)
  183. and isinstance(result.get("output_artifact"), str)
  184. ):
  185. response["output_artifact"] = result["output_artifact"]
  186. replay_body = (
  187. response
  188. if node.get("type") in {"rule.apply", "quality.check"}
  189. else None
  190. )
  191. if not finish_safely(
  192. claims.jti,
  193. status="success",
  194. commit_outcome=commit_outcome,
  195. safe_detail="task completed",
  196. replay_http_status=(200 if replay_body is not None else None),
  197. replay_body=replay_body,
  198. ):
  199. return jsonify({"error": "task ledger is unavailable"}), 503
  200. return jsonify(response)
  201. return app