| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223 |
- import json
- import logging
- from flask import g, jsonify, request
- from app import db
- from app.api.data_flow import bp
- from app.core.data_flow.dataflows import DataFlowService
- from app.core.graph.graph_operations import MyEncoder
- from app.models.result import failed, success
- from app.core.data_factory.n8n_client import N8nClient, N8nClientError
- from app.core.data_flow.workflow_activation import activate_version
- from app.core.data_flow.workflow_repository import create_version, list_versions
- from app.core.system.permissions import ACTIVATE_WORKFLOW, EDIT_GOVERNANCE, READ_GOVERNANCE, require_permissions
- logger = logging.getLogger(__name__)
- @bp.route("/<dataflow_uid>/workflow-versions", methods=["GET"])
- @require_permissions(READ_GOVERNANCE)
- def get_workflow_versions(dataflow_uid):
- try:
- return jsonify(success(list_versions(db.session, dataflow_uid)))
- except ValueError as exc:
- return jsonify(failed(str(exc), code=400)), 400
- @bp.route("/<dataflow_uid>/workflow-versions", methods=["POST"])
- @require_permissions(EDIT_GOVERNANCE)
- def add_workflow_version(dataflow_uid):
- body = request.get_json(silent=True) or {}
- n8n_workflow_id = str(body.get("n8n_workflow_id") or "").strip()
- if not n8n_workflow_id:
- return jsonify(failed("n8n_workflow_id 不能为空", code=400)), 400
- try:
- workflow = N8nClient().get_workflow(n8n_workflow_id)
- version_id = create_version(
- db.session,
- dataflow_uid=dataflow_uid,
- environment=body.get("environment", "development"),
- workflow=workflow,
- created_by=g.current_user["id"],
- )
- db.session.commit()
- return jsonify(success({"id": version_id}, "版本创建成功", code=201)), 201
- except (ValueError, N8nClientError) as exc:
- db.session.rollback()
- message = exc.message if isinstance(exc, N8nClientError) else str(exc)
- return jsonify(failed(message, code=400)), 400
- @bp.route("/<dataflow_uid>/workflow-versions/<version_id>/activate", methods=["POST"])
- @require_permissions(ACTIVATE_WORKFLOW)
- def activate_workflow_version(dataflow_uid, version_id):
- try:
- versions = list_versions(db.session, dataflow_uid)
- if not any(version["id"] == version_id for version in versions):
- return jsonify(failed("版本不存在", code=404)), 404
- result = activate_version(
- db.session,
- version_id=version_id,
- actor_id=g.current_user["id"],
- n8n_client=N8nClient(),
- )
- return jsonify(success(result, "版本激活成功"))
- except (ValueError, RuntimeError, N8nClientError) as exc:
- db.session.rollback()
- message = exc.message if isinstance(exc, N8nClientError) else str(exc)
- return jsonify(failed(message, code=409)), 409
- @bp.route("/get-dataflows-list", methods=["GET"])
- def get_dataflows():
- """获取数据流列表"""
- try:
- page = request.args.get("page", 1, type=int)
- page_size = request.args.get("page_size", 10, type=int)
- search = request.args.get("search", "")
- result = DataFlowService.get_dataflows(
- page=page,
- page_size=page_size,
- search=search,
- )
- res = success(result, "success")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except Exception as e:
- logger.error(f"获取数据流列表失败: {str(e)}")
- res = failed(f"获取数据流列表失败: {str(e)}")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- @bp.route("/get-dataflow/<int:dataflow_id>", methods=["GET"])
- def get_dataflow(dataflow_id):
- """根据ID获取数据流详情"""
- try:
- result = DataFlowService.get_dataflow_by_id(dataflow_id)
- if result:
- res = success(result, "success")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- else:
- res = failed("数据流不存在", code=404)
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except Exception as e:
- logger.error(f"获取数据流详情失败: {str(e)}")
- res = failed(f"获取数据流详情失败: {str(e)}")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- @bp.route("/add-dataflow", methods=["POST"])
- def create_dataflow():
- """创建新的数据流"""
- try:
- data = request.get_json()
- if not data:
- res = failed("请求数据不能为空", code=400)
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- result = DataFlowService.create_dataflow(data)
- res = success(result, "数据流创建成功")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except ValueError as ve:
- logger.error(f"创建数据流参数错误: {str(ve)}")
- res = failed(f"参数错误: {str(ve)}", code=400)
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except Exception as e:
- logger.error(f"创建数据流失败: {str(e)}")
- res = failed(f"创建数据流失败: {str(e)}")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- @bp.route("/update-dataflow/<int:dataflow_id>", methods=["PUT"])
- def update_dataflow(dataflow_id):
- """更新数据流"""
- try:
- data = request.get_json()
- if not data:
- res = failed("请求数据不能为空", code=400)
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- result = DataFlowService.update_dataflow(dataflow_id, data)
- if result:
- res = success(result, "数据流更新成功")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- else:
- res = failed("数据流不存在", code=404)
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except Exception as e:
- logger.error(f"更新数据流失败: {str(e)}")
- res = failed(f"更新数据流失败: {str(e)}")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- @bp.route("/delete-dataflow/<int:dataflow_id>", methods=["DELETE"])
- def delete_dataflow(dataflow_id):
- """删除数据流"""
- try:
- result = DataFlowService.delete_dataflow(dataflow_id)
- if result:
- res = success({}, "数据流删除成功")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- else:
- res = failed("数据流不存在", code=404)
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except Exception as e:
- logger.error(f"删除数据流失败: {str(e)}")
- res = failed(f"删除数据流失败: {str(e)}")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- @bp.route("/get-BD-list", methods=["GET"])
- def get_business_domain_list():
- """获取BusinessDomain节点列表"""
- try:
- logger.info("接收到获取BusinessDomain列表请求")
- # 调用服务层函数获取BusinessDomain列表
- bd_list = DataFlowService.get_business_domain_list()
- res = success(bd_list, "操作成功")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except Exception as e:
- logger.error(f"获取BusinessDomain列表失败: {str(e)}")
- res = failed(f"获取BusinessDomain列表失败: {str(e)}", 500, {})
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- @bp.route("/get-script/<int:dataflow_id>", methods=["GET"])
- def get_script(dataflow_id):
- """
- 获取 DataFlow 关联的脚本内容
- Args:
- dataflow_id: DataFlow 节点的 ID
- Returns:
- 包含脚本内容和元信息的 JSON 响应:
- - script_path: 脚本路径
- - script_content: 脚本内容
- - script_type: 脚本类型(python/javascript/sql等)
- - dataflow_id: DataFlow ID
- - dataflow_name: DataFlow 中文名称
- - dataflow_name_en: DataFlow 英文名称
- """
- try:
- logger.info(f"接收到获取脚本请求, DataFlow ID: {dataflow_id}")
- result = DataFlowService.get_script_content(dataflow_id)
- res = success(result, "获取脚本成功")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except ValueError as ve:
- logger.warning(f"获取脚本参数错误: {str(ve)}")
- res = failed(f"{str(ve)}", code=400)
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except FileNotFoundError as fe:
- logger.warning(f"脚本文件不存在: {str(fe)}")
- res = failed(f"脚本文件不存在: {str(fe)}", code=404)
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
- except Exception as e:
- logger.error(f"获取脚本失败: {str(e)}")
- res = failed(f"获取脚本失败: {str(e)}")
- return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
|