"""Minimal HTTP surface for the independent DataOps Runner.""" from flask import Flask, jsonify, request from app.runner.auth import TaskTokenInvalid from app.runner.nodes import NodeExecutionError def create_runner_app(*, verifier, ledger, registry): app = Flask("dataops-runner") app.config["MAX_CONTENT_LENGTH"] = 2 * 1024 * 1024 @app.get("/health") def health(): return jsonify({"status": "ok"}) def finish_safely(jti, **outcome): try: ledger.finish(jti, **outcome) return True except Exception: app.logger.error("runner task ledger unavailable") return False @app.post("/v1/tasks/execute") def execute_task(): payload = request.get_json(silent=True) if not isinstance(payload, dict): return jsonify({"error": "invalid task request"}), 400 node = payload.get("node") parameters = payload.get("parameters", {}) if not isinstance(node, dict) or not isinstance(parameters, dict): return jsonify({"error": "invalid task request"}), 400 try: claims = verifier.verify(payload.get("task_token"), node=node) except TaskTokenInvalid: return jsonify({"error": "task token is invalid"}), 401 binding = { "task_uid": claims.task_uid, "dataflow_uid": claims.dataflow_uid, "workflow_version": claims.workflow_version, "correlation_id": claims.correlation_id, "node_id": claims.node_id, "node_type": claims.node_type, "data_source_uid": node.get("data_source_uid"), "idempotency_key": (node.get("idempotency") or {}).get("key"), } try: claimed = ledger.claim( claims.jti, binding, expires_at=claims.expires_at, ) except Exception: app.logger.error("runner task ledger unavailable") return jsonify({"error": "task ledger is unavailable"}), 503 if not claimed: return jsonify({"error": "task token already consumed"}), 409 try: result = registry.execute( node, parameters, write_authorized=claims.write_authorized, ) except NodeExecutionError as exc: recorded = finish_safely( claims.jti, status=( "unknown" if exc.commit_outcome == "unknown" else "failed" ), commit_outcome=exc.commit_outcome, safe_detail=str(exc), ) if not recorded: return jsonify({"error": "task ledger is unavailable"}), 503 return jsonify({"error": str(exc)}), 400 except Exception: finish_safely( claims.jti, status="failed", safe_detail="task execution failed", ) app.logger.exception("runner task failed") return jsonify({"error": "task execution failed"}), 500 commit_outcome = ( result.get("commit_outcome", "not_applicable") if isinstance(result, dict) else "not_applicable" ) if not finish_safely( claims.jti, status="success", commit_outcome=commit_outcome, safe_detail="task completed", ): return jsonify({"error": "task ledger is unavailable"}), 503 return jsonify( { "task_uid": claims.task_uid, "correlation_id": claims.correlation_id, "result": result, } ) return app