api.py 15 KB

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