api.py 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349
  1. """Minimal HTTP surface for the independent DataOps Runner."""
  2. import hashlib
  3. import json
  4. import re
  5. from flask import Flask, jsonify, request
  6. from app.runner.auth import TaskTokenExpired, TaskTokenInvalid
  7. from app.runner.nodes import NodeExecutionError
  8. def create_runner_app(*, verifier, ledger, registry):
  9. app = Flask("dataops-runner")
  10. app.config["MAX_CONTENT_LENGTH"] = 2 * 1024 * 1024
  11. @app.get("/health")
  12. def health():
  13. return jsonify({"status": "ok"})
  14. def finish_safely(jti, **outcome):
  15. try:
  16. ledger.finish(jti, **outcome)
  17. return True
  18. except Exception:
  19. app.logger.error("runner task ledger unavailable")
  20. return False
  21. @app.post("/v1/tasks/execute")
  22. def execute_task():
  23. payload = request.get_json(silent=True)
  24. if not isinstance(payload, dict):
  25. return jsonify({"error": "invalid task request"}), 400
  26. node = payload.get("node")
  27. parameters = payload.get("parameters", {})
  28. if not isinstance(node, dict) or not isinstance(parameters, dict):
  29. return jsonify({"error": "invalid task request"}), 400
  30. replay_only = False
  31. try:
  32. claims = verifier.verify(payload.get("task_token"), node=node)
  33. except TaskTokenExpired:
  34. try:
  35. claims = verifier.verify_for_replay(
  36. payload.get("task_token"),
  37. node=node,
  38. )
  39. except TaskTokenInvalid:
  40. return jsonify({"error": "task token is invalid"}), 401
  41. replay_only = True
  42. except TaskTokenInvalid:
  43. return jsonify({"error": "task token is invalid"}), 401
  44. binding = {
  45. "task_uid": claims.task_uid,
  46. "dataflow_uid": claims.dataflow_uid,
  47. "deployment_id": claims.deployment_id,
  48. "environment": claims.environment,
  49. "workflow_version": claims.workflow_version,
  50. "correlation_id": claims.correlation_id,
  51. "node_id": claims.node_id,
  52. "node_type": claims.node_type,
  53. "data_source_uid": node.get("data_source_uid"),
  54. "idempotency_key": (node.get("idempotency") or {}).get("key"),
  55. }
  56. if replay_only:
  57. claimed = False
  58. else:
  59. try:
  60. claimed = ledger.claim(
  61. claims.jti,
  62. binding,
  63. expires_at=claims.expires_at,
  64. )
  65. except Exception:
  66. app.logger.error("runner task ledger unavailable")
  67. return jsonify({"error": "task ledger is unavailable"}), 503
  68. if not claimed:
  69. try:
  70. existing = ledger.get(claims.jti)
  71. except Exception:
  72. return jsonify({"error": "task ledger is unavailable"}), 503
  73. if (
  74. existing is None
  75. or any(
  76. existing.binding.get(key) != value
  77. for key, value in binding.items()
  78. )
  79. ):
  80. status = 401 if replay_only else 409
  81. return jsonify({"error": "task token already consumed"}), status
  82. if existing.status == "running":
  83. recovered = None
  84. if node.get("type") in {"rule.apply", "quality.check"}:
  85. try:
  86. recovered = registry.replay_task(
  87. node,
  88. correlation_id=claims.correlation_id,
  89. deployment_id=claims.deployment_id,
  90. task_jti=claims.jti,
  91. )
  92. except NodeExecutionError:
  93. return jsonify(
  94. {"error": "task ledger is unavailable"}
  95. ), 503
  96. recovered_state = (
  97. recovered.get("state")
  98. if isinstance(recovered, dict)
  99. else None
  100. )
  101. if isinstance(recovered, dict):
  102. expected_fields = {
  103. "running": {"state"},
  104. "missing": {"state"},
  105. "terminal": {
  106. "state",
  107. "status",
  108. "commit_outcome",
  109. "result",
  110. "result_digest",
  111. "evidence_digest",
  112. },
  113. }.get(recovered_state)
  114. if (
  115. expected_fields is None
  116. or set(recovered) != expected_fields
  117. ):
  118. return jsonify(
  119. {"error": "task ledger is unavailable"}
  120. ), 503
  121. recovered_result = (
  122. recovered.get("result")
  123. if recovered_state == "terminal"
  124. and recovered.get("status") == "success"
  125. else None
  126. )
  127. if (
  128. recovered_state == "terminal"
  129. and recovered.get("status") != "success"
  130. ):
  131. if not finish_safely(
  132. claims.jti,
  133. status="unknown",
  134. commit_outcome=str(
  135. recovered.get("commit_outcome")
  136. or "unknown"
  137. ),
  138. safe_detail=(
  139. "task execution evidence is terminal"
  140. ),
  141. ):
  142. return jsonify(
  143. {"error": "task ledger is unavailable"}
  144. ), 503
  145. return jsonify(
  146. {
  147. "error": (
  148. "task execution outcome is unknown"
  149. )
  150. }
  151. ), 409
  152. if recovered_state == "terminal":
  153. result_digest = str(
  154. recovered.get("result_digest") or ""
  155. )
  156. evidence_digest = str(
  157. recovered.get("evidence_digest") or ""
  158. )
  159. if (
  160. not isinstance(recovered_result, dict)
  161. or recovered.get("commit_outcome")
  162. != recovered_result.get("commit_outcome")
  163. or re.fullmatch(
  164. r"[0-9a-f]{64}",
  165. result_digest,
  166. )
  167. is None
  168. or re.fullmatch(
  169. r"[0-9a-f]{64}",
  170. evidence_digest,
  171. )
  172. is None
  173. ):
  174. return jsonify(
  175. {"error": "task ledger is unavailable"}
  176. ), 503
  177. encoded_result = json.dumps(
  178. recovered_result,
  179. sort_keys=True,
  180. separators=(",", ":"),
  181. ensure_ascii=False,
  182. ).encode("utf-8")
  183. if (
  184. hashlib.sha256(encoded_result).hexdigest()
  185. != result_digest
  186. ):
  187. return jsonify(
  188. {"error": "task ledger is unavailable"}
  189. ), 503
  190. replay_body = {
  191. "task_uid": claims.task_uid,
  192. "correlation_id": claims.correlation_id,
  193. "result": recovered_result,
  194. }
  195. output_artifact = recovered_result.get(
  196. "output_artifact"
  197. )
  198. if isinstance(output_artifact, str):
  199. replay_body["output_artifact"] = output_artifact
  200. if not finish_safely(
  201. claims.jti,
  202. status="success",
  203. commit_outcome=recovered_result.get(
  204. "commit_outcome",
  205. "not_applicable",
  206. ),
  207. safe_detail="task response recovered",
  208. replay_http_status=200,
  209. replay_body=replay_body,
  210. ):
  211. return jsonify(
  212. {"error": "task ledger is unavailable"}
  213. ), 503
  214. response = jsonify(replay_body)
  215. response.headers["X-Idempotent-Replay"] = "true"
  216. return response
  217. if recovered_state in {None, "missing"}:
  218. try:
  219. existing = ledger.reconcile_running(
  220. claims.jti,
  221. binding,
  222. )
  223. except Exception:
  224. return jsonify(
  225. {"error": "task ledger is unavailable"}
  226. ), 503
  227. if (
  228. existing is not None
  229. and existing.status == "unknown"
  230. ):
  231. return jsonify(
  232. {
  233. "error": (
  234. "task execution outcome is unknown"
  235. )
  236. }
  237. ), 409
  238. elif recovered_state != "running":
  239. return jsonify(
  240. {"error": "task ledger is unavailable"}
  241. ), 503
  242. if replay_only:
  243. return jsonify(
  244. {"error": "expired task is not replayable"}
  245. ), 401
  246. response = jsonify(
  247. {
  248. "task_uid": claims.task_uid,
  249. "correlation_id": claims.correlation_id,
  250. "status": "running",
  251. }
  252. )
  253. response.headers["Retry-After"] = "2"
  254. return response, 202
  255. if (
  256. existing.replay_body is None
  257. or existing.replay_http_status is None
  258. or existing.replay_digest is None
  259. ):
  260. return jsonify({"error": "task token already consumed"}), 409
  261. encoded = json.dumps(
  262. existing.replay_body,
  263. sort_keys=True,
  264. separators=(",", ":"),
  265. ensure_ascii=False,
  266. ).encode("utf-8")
  267. if (
  268. hashlib.sha256(encoded).hexdigest()
  269. != existing.replay_digest
  270. ):
  271. return jsonify({"error": "task ledger is unavailable"}), 503
  272. response = jsonify(existing.replay_body)
  273. response.headers["X-Idempotent-Replay"] = "true"
  274. return response, existing.replay_http_status
  275. try:
  276. result = registry.execute(
  277. node,
  278. parameters,
  279. write_authorized=claims.write_authorized,
  280. correlation_id=claims.correlation_id,
  281. dataflow_uid=claims.dataflow_uid,
  282. deployment_id=claims.deployment_id,
  283. environment=claims.environment,
  284. workflow_version=claims.workflow_version,
  285. node_id=claims.node_id,
  286. task_jti=claims.jti,
  287. )
  288. except NodeExecutionError as exc:
  289. recorded = finish_safely(
  290. claims.jti,
  291. status=(
  292. "unknown"
  293. if exc.commit_outcome == "unknown"
  294. else "failed"
  295. ),
  296. commit_outcome=exc.commit_outcome,
  297. safe_detail=str(exc),
  298. )
  299. if not recorded:
  300. return jsonify({"error": "task ledger is unavailable"}), 503
  301. return jsonify({"error": str(exc)}), 400
  302. except Exception:
  303. finish_safely(
  304. claims.jti,
  305. status="failed",
  306. safe_detail="task execution failed",
  307. )
  308. app.logger.exception("runner task failed")
  309. return jsonify({"error": "task execution failed"}), 500
  310. commit_outcome = (
  311. result.get("commit_outcome", "not_applicable")
  312. if isinstance(result, dict)
  313. else "not_applicable"
  314. )
  315. response = {
  316. "task_uid": claims.task_uid,
  317. "correlation_id": claims.correlation_id,
  318. "result": result,
  319. }
  320. if (
  321. isinstance(result, dict)
  322. and isinstance(result.get("output_artifact"), str)
  323. ):
  324. response["output_artifact"] = result["output_artifact"]
  325. replay_body = (
  326. response
  327. if node.get("type") in {"rule.apply", "quality.check"}
  328. else None
  329. )
  330. if not finish_safely(
  331. claims.jti,
  332. status="success",
  333. commit_outcome=commit_outcome,
  334. safe_detail="task completed",
  335. replay_http_status=(200 if replay_body is not None else None),
  336. replay_body=replay_body,
  337. ):
  338. return jsonify({"error": "task ledger is unavailable"}), 503
  339. return jsonify(response)
  340. return app