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("//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("//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("//workflow-versions//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/", 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/", 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/", 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/", 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)