api.py 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115
  1. """Minimal HTTP surface for the independent DataOps Runner."""
  2. from flask import Flask, jsonify, request
  3. from app.runner.auth import TaskTokenInvalid
  4. from app.runner.nodes import NodeExecutionError
  5. def create_runner_app(*, verifier, ledger, registry):
  6. app = Flask("dataops-runner")
  7. app.config["MAX_CONTENT_LENGTH"] = 2 * 1024 * 1024
  8. @app.get("/health")
  9. def health():
  10. return jsonify({"status": "ok"})
  11. def finish_safely(jti, **outcome):
  12. try:
  13. ledger.finish(jti, **outcome)
  14. return True
  15. except Exception:
  16. app.logger.error("runner task ledger unavailable")
  17. return False
  18. @app.post("/v1/tasks/execute")
  19. def execute_task():
  20. payload = request.get_json(silent=True)
  21. if not isinstance(payload, dict):
  22. return jsonify({"error": "invalid task request"}), 400
  23. node = payload.get("node")
  24. parameters = payload.get("parameters", {})
  25. if not isinstance(node, dict) or not isinstance(parameters, dict):
  26. return jsonify({"error": "invalid task request"}), 400
  27. try:
  28. claims = verifier.verify(payload.get("task_token"), node=node)
  29. except TaskTokenInvalid:
  30. return jsonify({"error": "task token is invalid"}), 401
  31. binding = {
  32. "task_uid": claims.task_uid,
  33. "dataflow_uid": claims.dataflow_uid,
  34. "workflow_version": claims.workflow_version,
  35. "correlation_id": claims.correlation_id,
  36. "node_id": claims.node_id,
  37. "node_type": claims.node_type,
  38. "data_source_uid": node.get("data_source_uid"),
  39. "idempotency_key": (node.get("idempotency") or {}).get("key"),
  40. }
  41. try:
  42. claimed = ledger.claim(
  43. claims.jti,
  44. binding,
  45. expires_at=claims.expires_at,
  46. )
  47. except Exception:
  48. app.logger.error("runner task ledger unavailable")
  49. return jsonify({"error": "task ledger is unavailable"}), 503
  50. if not claimed:
  51. return jsonify({"error": "task token already consumed"}), 409
  52. try:
  53. result = registry.execute(
  54. node,
  55. parameters,
  56. write_authorized=claims.write_authorized,
  57. correlation_id=claims.correlation_id,
  58. dataflow_uid=claims.dataflow_uid,
  59. workflow_version=claims.workflow_version,
  60. node_id=claims.node_id,
  61. )
  62. except NodeExecutionError as exc:
  63. recorded = finish_safely(
  64. claims.jti,
  65. status=(
  66. "unknown"
  67. if exc.commit_outcome == "unknown"
  68. else "failed"
  69. ),
  70. commit_outcome=exc.commit_outcome,
  71. safe_detail=str(exc),
  72. )
  73. if not recorded:
  74. return jsonify({"error": "task ledger is unavailable"}), 503
  75. return jsonify({"error": str(exc)}), 400
  76. except Exception:
  77. finish_safely(
  78. claims.jti,
  79. status="failed",
  80. safe_detail="task execution failed",
  81. )
  82. app.logger.exception("runner task failed")
  83. return jsonify({"error": "task execution failed"}), 500
  84. commit_outcome = (
  85. result.get("commit_outcome", "not_applicable")
  86. if isinstance(result, dict)
  87. else "not_applicable"
  88. )
  89. if not finish_safely(
  90. claims.jti,
  91. status="success",
  92. commit_outcome=commit_outcome,
  93. safe_detail="task completed",
  94. ):
  95. return jsonify({"error": "task ledger is unavailable"}), 503
  96. response = {
  97. "task_uid": claims.task_uid,
  98. "correlation_id": claims.correlation_id,
  99. "result": result,
  100. }
  101. if (
  102. isinstance(result, dict)
  103. and isinstance(result.get("output_artifact"), str)
  104. ):
  105. response["output_artifact"] = result["output_artifact"]
  106. return jsonify(response)
  107. return app