| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107 |
- """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
|