api.py 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290
  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 TaskTokenExpired, 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. replay_only = False
  30. try:
  31. claims = verifier.verify(payload.get("task_token"), node=node)
  32. except TaskTokenExpired:
  33. try:
  34. claims = verifier.verify_for_replay(
  35. payload.get("task_token"),
  36. node=node,
  37. )
  38. except TaskTokenInvalid:
  39. return jsonify({"error": "task token is invalid"}), 401
  40. replay_only = True
  41. except TaskTokenInvalid:
  42. return jsonify({"error": "task token is invalid"}), 401
  43. binding = {
  44. "task_uid": claims.task_uid,
  45. "dataflow_uid": claims.dataflow_uid,
  46. "deployment_id": claims.deployment_id,
  47. "environment": claims.environment,
  48. "workflow_version": claims.workflow_version,
  49. "correlation_id": claims.correlation_id,
  50. "node_id": claims.node_id,
  51. "node_type": claims.node_type,
  52. "data_source_uid": node.get("data_source_uid"),
  53. "idempotency_key": (node.get("idempotency") or {}).get("key"),
  54. }
  55. if replay_only:
  56. claimed = False
  57. else:
  58. try:
  59. claimed = ledger.claim(
  60. claims.jti,
  61. binding,
  62. expires_at=claims.expires_at,
  63. )
  64. except Exception:
  65. app.logger.error("runner task ledger unavailable")
  66. return jsonify({"error": "task ledger is unavailable"}), 503
  67. if not claimed:
  68. try:
  69. existing = ledger.get(claims.jti)
  70. except Exception:
  71. return jsonify({"error": "task ledger is unavailable"}), 503
  72. if (
  73. existing is None
  74. or any(
  75. existing.binding.get(key) != value
  76. for key, value in binding.items()
  77. )
  78. ):
  79. status = 401 if replay_only else 409
  80. return jsonify({"error": "task token already consumed"}), status
  81. if existing.status == "running":
  82. recovered = None
  83. if node.get("type") in {"rule.apply", "quality.check"}:
  84. try:
  85. recovered = registry.replay_task(
  86. node,
  87. correlation_id=claims.correlation_id,
  88. deployment_id=claims.deployment_id,
  89. task_jti=claims.jti,
  90. )
  91. except NodeExecutionError:
  92. return jsonify(
  93. {"error": "task ledger is unavailable"}
  94. ), 503
  95. recovered_state = (
  96. recovered.get("state")
  97. if isinstance(recovered, dict)
  98. else None
  99. )
  100. recovered_result = (
  101. recovered.get("result")
  102. if recovered_state == "terminal"
  103. and recovered.get("status") == "success"
  104. else recovered
  105. )
  106. if (
  107. recovered_state == "terminal"
  108. and recovered.get("status") != "success"
  109. ):
  110. if not finish_safely(
  111. claims.jti,
  112. status="unknown",
  113. commit_outcome=str(
  114. recovered.get("commit_outcome")
  115. or "unknown"
  116. ),
  117. safe_detail=(
  118. "task execution evidence is terminal"
  119. ),
  120. ):
  121. return jsonify(
  122. {"error": "task ledger is unavailable"}
  123. ), 503
  124. return jsonify(
  125. {
  126. "error": (
  127. "task execution outcome is unknown"
  128. )
  129. }
  130. ), 409
  131. if (
  132. isinstance(recovered_result, dict)
  133. and recovered_state != "missing"
  134. ):
  135. replay_body = {
  136. "task_uid": claims.task_uid,
  137. "correlation_id": claims.correlation_id,
  138. "result": recovered_result,
  139. }
  140. output_artifact = recovered_result.get(
  141. "output_artifact"
  142. )
  143. if isinstance(output_artifact, str):
  144. replay_body["output_artifact"] = output_artifact
  145. if not finish_safely(
  146. claims.jti,
  147. status="success",
  148. commit_outcome=recovered_result.get(
  149. "commit_outcome",
  150. "not_applicable",
  151. ),
  152. safe_detail="task response recovered",
  153. replay_http_status=200,
  154. replay_body=replay_body,
  155. ):
  156. return jsonify(
  157. {"error": "task ledger is unavailable"}
  158. ), 503
  159. response = jsonify(replay_body)
  160. response.headers["X-Idempotent-Replay"] = "true"
  161. return response
  162. if recovered_state in {None, "missing"}:
  163. try:
  164. existing = ledger.reconcile_running(
  165. claims.jti,
  166. binding,
  167. )
  168. except Exception:
  169. return jsonify(
  170. {"error": "task ledger is unavailable"}
  171. ), 503
  172. if (
  173. existing is not None
  174. and existing.status == "unknown"
  175. ):
  176. return jsonify(
  177. {
  178. "error": (
  179. "task execution outcome is unknown"
  180. )
  181. }
  182. ), 409
  183. if replay_only:
  184. return jsonify(
  185. {"error": "expired task is not replayable"}
  186. ), 401
  187. response = jsonify(
  188. {
  189. "task_uid": claims.task_uid,
  190. "correlation_id": claims.correlation_id,
  191. "status": "running",
  192. }
  193. )
  194. response.headers["Retry-After"] = "2"
  195. return response, 202
  196. if (
  197. existing.replay_body is None
  198. or existing.replay_http_status is None
  199. or existing.replay_digest is None
  200. ):
  201. return jsonify({"error": "task token already consumed"}), 409
  202. encoded = json.dumps(
  203. existing.replay_body,
  204. sort_keys=True,
  205. separators=(",", ":"),
  206. ensure_ascii=False,
  207. ).encode("utf-8")
  208. if (
  209. hashlib.sha256(encoded).hexdigest()
  210. != existing.replay_digest
  211. ):
  212. return jsonify({"error": "task ledger is unavailable"}), 503
  213. response = jsonify(existing.replay_body)
  214. response.headers["X-Idempotent-Replay"] = "true"
  215. return response, existing.replay_http_status
  216. try:
  217. result = registry.execute(
  218. node,
  219. parameters,
  220. write_authorized=claims.write_authorized,
  221. correlation_id=claims.correlation_id,
  222. dataflow_uid=claims.dataflow_uid,
  223. deployment_id=claims.deployment_id,
  224. environment=claims.environment,
  225. workflow_version=claims.workflow_version,
  226. node_id=claims.node_id,
  227. task_jti=claims.jti,
  228. )
  229. except NodeExecutionError as exc:
  230. recorded = finish_safely(
  231. claims.jti,
  232. status=(
  233. "unknown"
  234. if exc.commit_outcome == "unknown"
  235. else "failed"
  236. ),
  237. commit_outcome=exc.commit_outcome,
  238. safe_detail=str(exc),
  239. )
  240. if not recorded:
  241. return jsonify({"error": "task ledger is unavailable"}), 503
  242. return jsonify({"error": str(exc)}), 400
  243. except Exception:
  244. finish_safely(
  245. claims.jti,
  246. status="failed",
  247. safe_detail="task execution failed",
  248. )
  249. app.logger.exception("runner task failed")
  250. return jsonify({"error": "task execution failed"}), 500
  251. commit_outcome = (
  252. result.get("commit_outcome", "not_applicable")
  253. if isinstance(result, dict)
  254. else "not_applicable"
  255. )
  256. response = {
  257. "task_uid": claims.task_uid,
  258. "correlation_id": claims.correlation_id,
  259. "result": result,
  260. }
  261. if (
  262. isinstance(result, dict)
  263. and isinstance(result.get("output_artifact"), str)
  264. ):
  265. response["output_artifact"] = result["output_artifact"]
  266. replay_body = (
  267. response
  268. if node.get("type") in {"rule.apply", "quality.check"}
  269. else None
  270. )
  271. if not finish_safely(
  272. claims.jti,
  273. status="success",
  274. commit_outcome=commit_outcome,
  275. safe_detail="task completed",
  276. replay_http_status=(200 if replay_body is not None else None),
  277. replay_body=replay_body,
  278. ):
  279. return jsonify({"error": "task ledger is unavailable"}), 503
  280. return jsonify(response)
  281. return app