api.py 3.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107
  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. )
  58. except NodeExecutionError as exc:
  59. recorded = finish_safely(
  60. claims.jti,
  61. status=(
  62. "unknown"
  63. if exc.commit_outcome == "unknown"
  64. else "failed"
  65. ),
  66. commit_outcome=exc.commit_outcome,
  67. safe_detail=str(exc),
  68. )
  69. if not recorded:
  70. return jsonify({"error": "task ledger is unavailable"}), 503
  71. return jsonify({"error": str(exc)}), 400
  72. except Exception:
  73. finish_safely(
  74. claims.jti,
  75. status="failed",
  76. safe_detail="task execution failed",
  77. )
  78. app.logger.exception("runner task failed")
  79. return jsonify({"error": "task execution failed"}), 500
  80. commit_outcome = (
  81. result.get("commit_outcome", "not_applicable")
  82. if isinstance(result, dict)
  83. else "not_applicable"
  84. )
  85. if not finish_safely(
  86. claims.jti,
  87. status="success",
  88. commit_outcome=commit_outcome,
  89. safe_detail="task completed",
  90. ):
  91. return jsonify({"error": "task ledger is unavailable"}), 503
  92. return jsonify(
  93. {
  94. "task_uid": claims.task_uid,
  95. "correlation_id": claims.correlation_id,
  96. "result": result,
  97. }
  98. )
  99. return app