| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349 |
- """Minimal HTTP surface for the independent DataOps Runner."""
- import hashlib
- import json
- import re
- from flask import Flask, jsonify, request
- from app.runner.auth import TaskTokenExpired, 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
- replay_only = False
- try:
- claims = verifier.verify(payload.get("task_token"), node=node)
- except TaskTokenExpired:
- try:
- claims = verifier.verify_for_replay(
- payload.get("task_token"),
- node=node,
- )
- except TaskTokenInvalid:
- return jsonify({"error": "task token is invalid"}), 401
- replay_only = True
- except TaskTokenInvalid:
- return jsonify({"error": "task token is invalid"}), 401
- binding = {
- "task_uid": claims.task_uid,
- "dataflow_uid": claims.dataflow_uid,
- "deployment_id": claims.deployment_id,
- "environment": claims.environment,
- "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"),
- }
- if replay_only:
- claimed = False
- else:
- 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:
- try:
- existing = ledger.get(claims.jti)
- except Exception:
- return jsonify({"error": "task ledger is unavailable"}), 503
- if (
- existing is None
- or any(
- existing.binding.get(key) != value
- for key, value in binding.items()
- )
- ):
- status = 401 if replay_only else 409
- return jsonify({"error": "task token already consumed"}), status
- if existing.status == "running":
- recovered = None
- if node.get("type") in {"rule.apply", "quality.check"}:
- try:
- recovered = registry.replay_task(
- node,
- correlation_id=claims.correlation_id,
- deployment_id=claims.deployment_id,
- task_jti=claims.jti,
- )
- except NodeExecutionError:
- return jsonify(
- {"error": "task ledger is unavailable"}
- ), 503
- recovered_state = (
- recovered.get("state")
- if isinstance(recovered, dict)
- else None
- )
- if isinstance(recovered, dict):
- expected_fields = {
- "running": {"state"},
- "missing": {"state"},
- "terminal": {
- "state",
- "status",
- "commit_outcome",
- "result",
- "result_digest",
- "evidence_digest",
- },
- }.get(recovered_state)
- if (
- expected_fields is None
- or set(recovered) != expected_fields
- ):
- return jsonify(
- {"error": "task ledger is unavailable"}
- ), 503
- recovered_result = (
- recovered.get("result")
- if recovered_state == "terminal"
- and recovered.get("status") == "success"
- else None
- )
- if (
- recovered_state == "terminal"
- and recovered.get("status") != "success"
- ):
- if not finish_safely(
- claims.jti,
- status="unknown",
- commit_outcome=str(
- recovered.get("commit_outcome")
- or "unknown"
- ),
- safe_detail=(
- "task execution evidence is terminal"
- ),
- ):
- return jsonify(
- {"error": "task ledger is unavailable"}
- ), 503
- return jsonify(
- {
- "error": (
- "task execution outcome is unknown"
- )
- }
- ), 409
- if recovered_state == "terminal":
- result_digest = str(
- recovered.get("result_digest") or ""
- )
- evidence_digest = str(
- recovered.get("evidence_digest") or ""
- )
- if (
- not isinstance(recovered_result, dict)
- or recovered.get("commit_outcome")
- != recovered_result.get("commit_outcome")
- or re.fullmatch(
- r"[0-9a-f]{64}",
- result_digest,
- )
- is None
- or re.fullmatch(
- r"[0-9a-f]{64}",
- evidence_digest,
- )
- is None
- ):
- return jsonify(
- {"error": "task ledger is unavailable"}
- ), 503
- encoded_result = json.dumps(
- recovered_result,
- sort_keys=True,
- separators=(",", ":"),
- ensure_ascii=False,
- ).encode("utf-8")
- if (
- hashlib.sha256(encoded_result).hexdigest()
- != result_digest
- ):
- return jsonify(
- {"error": "task ledger is unavailable"}
- ), 503
- replay_body = {
- "task_uid": claims.task_uid,
- "correlation_id": claims.correlation_id,
- "result": recovered_result,
- }
- output_artifact = recovered_result.get(
- "output_artifact"
- )
- if isinstance(output_artifact, str):
- replay_body["output_artifact"] = output_artifact
- if not finish_safely(
- claims.jti,
- status="success",
- commit_outcome=recovered_result.get(
- "commit_outcome",
- "not_applicable",
- ),
- safe_detail="task response recovered",
- replay_http_status=200,
- replay_body=replay_body,
- ):
- return jsonify(
- {"error": "task ledger is unavailable"}
- ), 503
- response = jsonify(replay_body)
- response.headers["X-Idempotent-Replay"] = "true"
- return response
- if recovered_state in {None, "missing"}:
- try:
- existing = ledger.reconcile_running(
- claims.jti,
- binding,
- )
- except Exception:
- return jsonify(
- {"error": "task ledger is unavailable"}
- ), 503
- if (
- existing is not None
- and existing.status == "unknown"
- ):
- return jsonify(
- {
- "error": (
- "task execution outcome is unknown"
- )
- }
- ), 409
- elif recovered_state != "running":
- return jsonify(
- {"error": "task ledger is unavailable"}
- ), 503
- if replay_only:
- return jsonify(
- {"error": "expired task is not replayable"}
- ), 401
- response = jsonify(
- {
- "task_uid": claims.task_uid,
- "correlation_id": claims.correlation_id,
- "status": "running",
- }
- )
- response.headers["Retry-After"] = "2"
- return response, 202
- if (
- existing.replay_body is None
- or existing.replay_http_status is None
- or existing.replay_digest is None
- ):
- return jsonify({"error": "task token already consumed"}), 409
- encoded = json.dumps(
- existing.replay_body,
- sort_keys=True,
- separators=(",", ":"),
- ensure_ascii=False,
- ).encode("utf-8")
- if (
- hashlib.sha256(encoded).hexdigest()
- != existing.replay_digest
- ):
- return jsonify({"error": "task ledger is unavailable"}), 503
- response = jsonify(existing.replay_body)
- response.headers["X-Idempotent-Replay"] = "true"
- return response, existing.replay_http_status
- try:
- result = registry.execute(
- node,
- parameters,
- write_authorized=claims.write_authorized,
- correlation_id=claims.correlation_id,
- dataflow_uid=claims.dataflow_uid,
- deployment_id=claims.deployment_id,
- environment=claims.environment,
- workflow_version=claims.workflow_version,
- node_id=claims.node_id,
- task_jti=claims.jti,
- )
- 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"
- )
- response = {
- "task_uid": claims.task_uid,
- "correlation_id": claims.correlation_id,
- "result": result,
- }
- if (
- isinstance(result, dict)
- and isinstance(result.get("output_artifact"), str)
- ):
- response["output_artifact"] = result["output_artifact"]
- replay_body = (
- response
- if node.get("type") in {"rule.apply", "quality.check"}
- else None
- )
- if not finish_safely(
- claims.jti,
- status="success",
- commit_outcome=commit_outcome,
- safe_detail="task completed",
- replay_http_status=(200 if replay_body is not None else None),
- replay_body=replay_body,
- ):
- return jsonify({"error": "task ledger is unavailable"}), 503
- return jsonify(response)
- return app
|