routes.py 9.1 KB

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