routes.py 8.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223
  1. import json
  2. import logging
  3. from flask import g, jsonify, request
  4. from app import db
  5. from app.api.data_flow import bp
  6. from app.core.data_flow.dataflows import DataFlowService
  7. from app.core.graph.graph_operations import MyEncoder
  8. from app.models.result import failed, success
  9. from app.core.data_factory.n8n_client import N8nClient, N8nClientError
  10. from app.core.data_flow.workflow_activation import activate_version
  11. from app.core.data_flow.workflow_repository import create_version, list_versions
  12. from app.core.system.permissions import ACTIVATE_WORKFLOW, EDIT_GOVERNANCE, READ_GOVERNANCE, require_permissions
  13. logger = logging.getLogger(__name__)
  14. @bp.route("/<dataflow_uid>/workflow-versions", methods=["GET"])
  15. @require_permissions(READ_GOVERNANCE)
  16. def get_workflow_versions(dataflow_uid):
  17. try:
  18. return jsonify(success(list_versions(db.session, dataflow_uid)))
  19. except ValueError as exc:
  20. return jsonify(failed(str(exc), code=400)), 400
  21. @bp.route("/<dataflow_uid>/workflow-versions", methods=["POST"])
  22. @require_permissions(EDIT_GOVERNANCE)
  23. def add_workflow_version(dataflow_uid):
  24. body = request.get_json(silent=True) or {}
  25. n8n_workflow_id = str(body.get("n8n_workflow_id") or "").strip()
  26. if not n8n_workflow_id:
  27. return jsonify(failed("n8n_workflow_id 不能为空", code=400)), 400
  28. try:
  29. workflow = N8nClient().get_workflow(n8n_workflow_id)
  30. version_id = create_version(
  31. db.session,
  32. dataflow_uid=dataflow_uid,
  33. environment=body.get("environment", "development"),
  34. workflow=workflow,
  35. created_by=g.current_user["id"],
  36. )
  37. db.session.commit()
  38. return jsonify(success({"id": version_id}, "版本创建成功", code=201)), 201
  39. except (ValueError, N8nClientError) as exc:
  40. db.session.rollback()
  41. message = exc.message if isinstance(exc, N8nClientError) else str(exc)
  42. return jsonify(failed(message, code=400)), 400
  43. @bp.route("/<dataflow_uid>/workflow-versions/<version_id>/activate", methods=["POST"])
  44. @require_permissions(ACTIVATE_WORKFLOW)
  45. def activate_workflow_version(dataflow_uid, version_id):
  46. try:
  47. versions = list_versions(db.session, dataflow_uid)
  48. if not any(version["id"] == version_id for version in versions):
  49. return jsonify(failed("版本不存在", code=404)), 404
  50. result = activate_version(
  51. db.session,
  52. version_id=version_id,
  53. actor_id=g.current_user["id"],
  54. n8n_client=N8nClient(),
  55. )
  56. return jsonify(success(result, "版本激活成功"))
  57. except (ValueError, RuntimeError, N8nClientError) as exc:
  58. db.session.rollback()
  59. message = exc.message if isinstance(exc, N8nClientError) else str(exc)
  60. return jsonify(failed(message, code=409)), 409
  61. @bp.route("/get-dataflows-list", methods=["GET"])
  62. def get_dataflows():
  63. """获取数据流列表"""
  64. try:
  65. page = request.args.get("page", 1, type=int)
  66. page_size = request.args.get("page_size", 10, type=int)
  67. search = request.args.get("search", "")
  68. result = DataFlowService.get_dataflows(
  69. page=page,
  70. page_size=page_size,
  71. search=search,
  72. )
  73. res = success(result, "success")
  74. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  75. except Exception as e:
  76. logger.error(f"获取数据流列表失败: {str(e)}")
  77. res = failed(f"获取数据流列表失败: {str(e)}")
  78. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  79. @bp.route("/get-dataflow/<int:dataflow_id>", methods=["GET"])
  80. def get_dataflow(dataflow_id):
  81. """根据ID获取数据流详情"""
  82. try:
  83. result = DataFlowService.get_dataflow_by_id(dataflow_id)
  84. if result:
  85. res = success(result, "success")
  86. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  87. else:
  88. res = failed("数据流不存在", code=404)
  89. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  90. except Exception as e:
  91. logger.error(f"获取数据流详情失败: {str(e)}")
  92. res = failed(f"获取数据流详情失败: {str(e)}")
  93. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  94. @bp.route("/add-dataflow", methods=["POST"])
  95. def create_dataflow():
  96. """创建新的数据流"""
  97. try:
  98. data = request.get_json()
  99. if not data:
  100. res = failed("请求数据不能为空", code=400)
  101. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  102. result = DataFlowService.create_dataflow(data)
  103. res = success(result, "数据流创建成功")
  104. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  105. except ValueError as ve:
  106. logger.error(f"创建数据流参数错误: {str(ve)}")
  107. res = failed(f"参数错误: {str(ve)}", code=400)
  108. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  109. except Exception as e:
  110. logger.error(f"创建数据流失败: {str(e)}")
  111. res = failed(f"创建数据流失败: {str(e)}")
  112. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  113. @bp.route("/update-dataflow/<int:dataflow_id>", methods=["PUT"])
  114. def update_dataflow(dataflow_id):
  115. """更新数据流"""
  116. try:
  117. data = request.get_json()
  118. if not data:
  119. res = failed("请求数据不能为空", code=400)
  120. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  121. result = DataFlowService.update_dataflow(dataflow_id, data)
  122. if result:
  123. res = success(result, "数据流更新成功")
  124. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  125. else:
  126. res = failed("数据流不存在", code=404)
  127. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  128. except Exception as e:
  129. logger.error(f"更新数据流失败: {str(e)}")
  130. res = failed(f"更新数据流失败: {str(e)}")
  131. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  132. @bp.route("/delete-dataflow/<int:dataflow_id>", methods=["DELETE"])
  133. def delete_dataflow(dataflow_id):
  134. """删除数据流"""
  135. try:
  136. result = DataFlowService.delete_dataflow(dataflow_id)
  137. if result:
  138. res = success({}, "数据流删除成功")
  139. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  140. else:
  141. res = failed("数据流不存在", code=404)
  142. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  143. except Exception as e:
  144. logger.error(f"删除数据流失败: {str(e)}")
  145. res = failed(f"删除数据流失败: {str(e)}")
  146. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  147. @bp.route("/get-BD-list", methods=["GET"])
  148. def get_business_domain_list():
  149. """获取BusinessDomain节点列表"""
  150. try:
  151. logger.info("接收到获取BusinessDomain列表请求")
  152. # 调用服务层函数获取BusinessDomain列表
  153. bd_list = DataFlowService.get_business_domain_list()
  154. res = success(bd_list, "操作成功")
  155. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  156. except Exception as e:
  157. logger.error(f"获取BusinessDomain列表失败: {str(e)}")
  158. res = failed(f"获取BusinessDomain列表失败: {str(e)}", 500, {})
  159. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  160. @bp.route("/get-script/<int:dataflow_id>", methods=["GET"])
  161. def get_script(dataflow_id):
  162. """
  163. 获取 DataFlow 关联的脚本内容
  164. Args:
  165. dataflow_id: DataFlow 节点的 ID
  166. Returns:
  167. 包含脚本内容和元信息的 JSON 响应:
  168. - script_path: 脚本路径
  169. - script_content: 脚本内容
  170. - script_type: 脚本类型(python/javascript/sql等)
  171. - dataflow_id: DataFlow ID
  172. - dataflow_name: DataFlow 中文名称
  173. - dataflow_name_en: DataFlow 英文名称
  174. """
  175. try:
  176. logger.info(f"接收到获取脚本请求, DataFlow ID: {dataflow_id}")
  177. result = DataFlowService.get_script_content(dataflow_id)
  178. res = success(result, "获取脚本成功")
  179. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  180. except ValueError as ve:
  181. logger.warning(f"获取脚本参数错误: {str(ve)}")
  182. res = failed(f"{str(ve)}", code=400)
  183. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  184. except FileNotFoundError as fe:
  185. logger.warning(f"脚本文件不存在: {str(fe)}")
  186. res = failed(f"脚本文件不存在: {str(fe)}", code=404)
  187. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)
  188. except Exception as e:
  189. logger.error(f"获取脚本失败: {str(e)}")
  190. res = failed(f"获取脚本失败: {str(e)}")
  191. return json.dumps(res, ensure_ascii=False, cls=MyEncoder)