| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448 |
- import json
- import time
- import os
- import yaml
- from passlib.context import CryptContext
- from Log import logger
- from app.config.agent_base_url import RG_APP_TOKEN_LIST, RG_APP_NEW_TOKEN, DF_CHAT_API_KEY
- # from app.api import pwd_context
- from app.config.const import DIFY, ENV_CONF_PATH, RAGFLOW, smart_server, chat_server, workflow_server, TMP_DICT, \
- rg_api_token, Dialog_STATSU_ON, SYSTEM_ID
- from app.models import MenuCapacityModel, WebMenuModel, GroupModel, RoleModel, DialogModel, UserModel, UserAppModel, \
- cipher_suite, UserTokenModel, ApiTokenModel, ComplexChatModel, SystemDataModel
- from app.service.auth import UserAppDao
- from app.service.bisheng import BishengService
- from app.service.difyService import DifyService
- from app.service.ragflow import RagflowService
- from app.service.service_token import get_new_token
- from app.service.v2.app_driver.chat_data import ChatBaseApply
- from app.service.v2.app_register import AppRegisterDao
- from app.config.config import settings
- from app.service.v2.chat import get_app_token
- from app.utils.common import get_machine_id
- from app.utils.password_handle import generate_password, password_encrypted, password_decrypted
- pwd_context = CryptContext(schemes=["bcrypt"], deprecated="auto")
- async def dialog_menu_sync(db):
- menu_list = []
- with open(os.path.join(ENV_CONF_PATH, "menu_conf.json") , 'r', encoding='utf-8') as file:
- # 加载JSON数据
- data = json.load(file)
- menu_list = data.get("data", [])
- db.query(WebMenuModel).delete()
- db.query(MenuCapacityModel).delete()
- db.commit()
- for menu in menu_list:
- # print(menu)
- dialog = menu.pop("dialog", [])
- for i in dialog:
- capacity = MenuCapacityModel(menu_id=menu["id"], capacity_id=i["id"], capacity_type=i["agentType"],
- chat_id=i["id"] if not i["chat_id"] else i["chat_id"],
- chat_type=i["chat_type"])
- db.add(capacity)
- menu_obj = WebMenuModel(**menu)
- db.add(menu_obj)
- db.commit()
- async def create_menu_sync(db):
- # json_file_path = "env_conf/menu_conf.json.template"
- json_file_path = os.path.join(ENV_CONF_PATH, "menu_conf.json.template")
- with open(json_file_path, 'r', encoding='utf-8') as file:
- json_data = json.load(file).get("data", [])
- # for menu in json_data:
- # menu['dialog'].clear()
- dialogs = db.query(DialogModel).all()
- dialog_dict = {}
- for dialog in dialogs:
- if dialog.name not in dialog_dict:
- dialog_dict[dialog.name] = []
- dialog_dict[dialog.name].append({
- 'id': dialog.id,
- 'chat_id': dialog.id,
- 'chat_type': '',
- 'agentType': dialog.dialog_type
- })
- for menu in json_data:
- # if menu['title'] in dialog_dict:
- # for dialog in dialog_dict[menu['title']]:
- # new_dialog_item = {
- # 'id': dialog.id,
- # 'chat_id': dialog.id,
- # 'chat_type': '',
- # 'agentType': dialog.dialog_type
- # }
- menu['dialog']= dialog_dict.get(menu['title'], [])
- json_data = {"data": json_data}
- new_file_name = f"menu_conf.json.template"
- new_file_path = os.path.join(os.path.dirname(json_file_path), new_file_name)
- with open(new_file_path, 'w', encoding='utf-8') as new_file:
- json.dump(json_data, new_file, ensure_ascii=False, indent=4)
- return {
- "file_name": new_file_name,
- "json_data": json_data
- }
- async def default_group_sync(db):
- group = db.query(GroupModel).filter_by(group_type=2).first()
- if not group:
- logger.error("未初始默认组, 开始初始化!")
- try:
- group = GroupModel(group_name="默认用户组", group_description="默认组", group_type=2)
- db.add(group)
- db.commit()
- except Exception as e:
- logger.error(e)
- async def default_role_sync(db):
- role = db.query(RoleModel).filter_by(role_type=2).first()
- if not role:
- logger.error("未初始默认角色, 开始初始化!")
- try:
- group = RoleModel(id="morenjuese1234567890", name="默认角色", description="默认角色", role_type=2)
- db.add(group)
- db.commit()
- except Exception as e:
- logger.error(e)
- # async def app_register_sync(db):
- # app_dict = {}
- # with open(os.path.join(ENV_CONF_PATH, "app_register_conf.json"), 'r', encoding='utf-8') as file:
- # # 加载JSON数据
- # app_dict = json.load(file)
- # try:
- # for app_id, status in app_dict.items():
- # AppRegisterDao(db).update_and_insert_app(app_id, status)
- # except Exception as e:
- # logger.error(e)
- async def basic_agent_sync(db):
- agent_list = []
- complex_list = []
- with open(os.path.join(ENV_CONF_PATH, "default_agent_conf.json"), 'r', encoding='utf-8') as file:
- # 加载JSON数据
- agent_dict = json.load(file)
- agent_list = agent_dict.get("basic", [])
- complex_list = agent_dict.get("complex", [])
- user = db.query(UserModel).filter_by(permission="admin").first()
- for agent in agent_list:
- dialog = db.query(DialogModel).filter(DialogModel.id == agent["id"]).first()
- if dialog:
- try:
- dialog.name = agent["name"]
- dialog.description = agent["description"]
- dialog.icon = agent["icon"]
- dialog.mode = agent["mode"]
- dialog.parameters = json.dumps(agent["parameters"])
- db.commit()
- except Exception as e:
- logger.error(e)
- else:
- try:
- dialog = DialogModel(id=agent["id"], name=agent["name"], description=agent["description"],
- icon=agent["icon"], tenant_id=user.id if user else "", dialog_type=agent["dialogType"], mode=agent["mode"],parameters = json.dumps(agent["parameters"]))
- db.add(dialog)
- db.commit()
- db.refresh(dialog)
- except Exception as e:
- print(e)
- db.rollback()
- now_complex_list = []
- for agent in complex_list:
- now_complex_list.append(agent["id"])
- dialog = db.query(ComplexChatModel).filter(ComplexChatModel.id == agent["id"]).first()
- if dialog:
- try:
- dialog.name = agent["name"]
- dialog.description = agent["description"]
- dialog.icon = agent["icon"]
- dialog.mode = agent["mode"]
- dialog.chat_mode = agent["chat_mode"]
- dialog.status = Dialog_STATSU_ON
- # dialog.parameters = json.dumps(agent["parameters"])
- db.commit()
- except Exception as e:
- logger.error(e)
- else:
- try:
- dialog = ComplexChatModel(id=agent["id"], name=agent["name"], description=agent["description"],
- icon=agent["icon"], tenant_id=user.id if user else "", dialog_type=agent["dialogType"], mode=agent["mode"],chat_mode = agent["chat_mode"])
- db.add(dialog)
- db.commit()
- db.refresh(dialog)
- except Exception as e:
- print(e)
- db.rollback()
- for i in db.query(ComplexChatModel).filter(ComplexChatModel.status == "1").all():
- if i.id not in now_complex_list:
- try:
- db.query(ComplexChatModel).filter(ComplexChatModel.id==i.id).update(({"status": "0"}))
- db.commit()
- except:
- ...
- async def user_update_app(userid, db):
- user = db.query(UserModel).filter(UserModel.id == userid).first()
- if not user:
- raise Exception("User id not found")
- app_register = AppRegisterDao(db).get_apps()
- register_dict = {}
- token = ""
- app_password = await generate_password(10)
- crypt_password = UserAppModel.encrypted_password(app_password)
- for app in app_register:
- if app["id"] == 'ragflow_app':
- user_rag_app = db.query(UserAppModel).filter(UserAppModel.user_id == userid,
- UserAppModel.app_type == 'ragflow_app').all()
- if not user_rag_app:
- service = RagflowService(settings.fwr_base_url)
- register_info = await register_app(service, app["id"], app_password, token)
- if register_info:
- register_dict[app["id"]] = register_info
- app_name = register_info.get("name")
- app_id = register_info.get("id")
- app_email = register_info.get("email")
- await save_db(db, app_name, crypt_password, app_email, user.id, app_id, "ragflow_app")
- elif app["id"] == 'bisheng_app':
- user_bs_app = db.query(UserAppModel).filter(UserAppModel.user_id == userid,
- UserAppModel.app_type == 'bisheng_app').all()
- if not user_bs_app:
- service = BishengService(settings.sgb_base_url)
- register_info = await register_app(service, app["id"], app_password, token)
- if register_info:
- register_dict[app["id"]] = register_info
- app_name = register_info.get("name")
- app_id = register_info.get("id")
- app_email = register_info.get("email")
- await save_db(db, app_name, crypt_password, app_email, user.id, app_id, "bisheng_app")
- elif app["id"] == 'dify_app':
- user_df_app = db.query(UserAppModel).filter(UserAppModel.user_id == userid,
- UserAppModel.app_type == 'dify_app').all()
- if not user_df_app:
- admin_user = db.query(UserModel).filter(UserModel.permission == "admin").first()
- token = await get_new_token(db, admin_user.id, DIFY)
- if not token:
- print("用户注册获取dftoken失败!")
- service = DifyService(settings.dify_base_url)
- register_info = await register_app(service, app["id"], app_password, token)
- if register_info:
- register_dict[app["id"]] = register_info
- app_name = register_info.get("name")
- app_id = register_info.get("id")
- app_email = register_info.get("email")
- await save_db(db, app_name, crypt_password, app_email, user.id, app_id, "dify_app")
- else:
- raise Exception("未知注册应用---")
- async def register_app(service, app_id, app_password, token):
- name = app_id + str(int(time.time()))
- try:
- register_info = await service.register(name, app_password, token)
- return {"id": register_info.get("id"), "name": name, "email": register_info.get("email")}
- except Exception as e:
- print(f"Failed to register with {app_id}: {str(e)}")
- return None
- async def save_db(db, username, password, email, user_id, app_id, app_type):
- user_app_dao = UserAppDao(db)
- user_id = await user_app_dao.insert_user_app_data(username, password, email, user_id, app_id, app_type)
- if not user_id:
- raise Exception("Failed to register with app")
- print({"msg": "User registered successfully", "userFlag": user_id})
- async def admin_account_sync(db):
- try:
- config = {}
- app_dict = {}
- # tmp_dict = {chat_server:RAGFLOW, workflow_server:DIFY}
- now_account =[]
- with open(os.path.join(ENV_CONF_PATH, "admin.yaml"), 'r', encoding='utf-8') as file:
- # 加载JSON数据
- config = yaml.safe_load(file)
- account_list = db.query(UserTokenModel).all()
- for account in account_list:
- if account.id in config:
- if account.account != config[account.id]["account"] or account.password != config[account.id]["password"]:
- db.query(UserTokenModel).filter_by(id=account.id).update({"account": config[account.id]["account"],
- "password": config[account.id]["password"],
- "access_token": ""
- })
- now_account.append(account.id)
- else:
- db.query(UserTokenModel).filter_by(id=account.id).delete()
- for k, v in config.items():
- if k in TMP_DICT:
- app_dict[TMP_DICT[k]] = v.get("id")
- if k == smart_server:
- db_user = db.query(UserModel).filter(UserModel.username == config["smart_server"]["account"]).first()
- if db_user:
- print("admin_user_sync: 用户已经存在!")
- continue
- hashed_password = pwd_context.hash(await password_decrypted(config["smart_server"]["password"])) # config["smart_server"]["password"]
- user_model = UserModel(username=config["smart_server"]["account"], hashed_password=hashed_password,
- email="",
- phone="", login_name="", sync_flag="", creator=0, permission="admin")
- db.add(user_model)
- # db.commit()
- # db.refresh(user_model)
- else:
- if k not in now_account:
- new_account = UserTokenModel(id=k, account=v["account"], password=v["password"])
- db.add(new_account)
- db.commit()
- # with open(os.path.join(ENV_CONF_PATH, "app_register_conf.json"), 'r', encoding='utf-8') as file:
- # # 加载JSON数据
- # app_dict = json.load(file)
- try:
- for app_id, name in app_dict.items():
- AppRegisterDao(db).update_and_insert_app(app_id, 1, name)
- except Exception as e:
- logger.error(e)
- except Exception as e:
- print(e)
- db.rollback()
- async def admin_user_sync(db):
- try:
- config = {}
- with open(os.path.join(ENV_CONF_PATH, "admin.yaml"), 'r', encoding='utf-8') as file:
- # 加载JSON数据
- config = yaml.safe_load(file)
- # print(config)
- db_user = db.query(UserModel).filter(UserModel.username == config["smart_server"]["account"]).first()
- if db_user:
- print("admin_user_sync: 用户已经存在!")
- return
- # register_dict = {}
- #
- # for app in [RAGFLOW, DIFY]:
- # register_dict[app] = {"id": config[app].get("id", "123"), "name": config[app]["account"],
- # "pwd":config[app]["password"],
- # "email": config[app]["account"]}
- # 存储用户信息
- hashed_password = pwd_context.hash(config["smart_server"]["password"])
- user_model = UserModel(username=config["smart_server"]["account"], hashed_password=hashed_password, email="",
- phone="", login_name="", sync_flag="", creator=0, permission="admin")
- db.add(user_model)
- db.commit()
- db.refresh(user_model)
- # u_id = user_model.id
- # user_app_dao = UserAppDao(db)
- # for k, v in register_dict.items():
- # await user_app_dao.update_and_insert_data(v.get("name"), user_model.encrypted_password(v.get("pwd")), v.get("email"), u_id, str(v.get("id")), k)
- except Exception as e:
- print(e)
- db.rollback()
- async def sync_rg_api_token(db):
- token = ""
- try:
- app_token = db.query(ApiTokenModel).filter_by(app_id=rg_api_token).first()
- if app_token:
- print("rg_api_token: 已经存在!")
- return
- user_token = db.query(UserTokenModel).filter(UserTokenModel.id == chat_server).first()
- chat = ChatBaseApply()
- token_list_url = f"{settings.fwr_base_url}{RG_APP_TOKEN_LIST}"
- token_list = await chat.chat_get(token_list_url, {}, await chat.get_chat_headers(user_token.access_token))
- if token_list and token_list.get("code") == 0:
- if len(token_list.get("data", [])) == 0:
- print("rg_api_token: 创建成功!")
- new_token_url = f"{settings.fwr_base_url}{RG_APP_NEW_TOKEN}"
- new_token = await chat.chat_post(new_token_url, {}, await chat.get_chat_headers(user_token.access_token))
- if new_token and new_token.get("code") == 0:
- token = new_token.get("data", {}).get("token")
- else:
- token = token_list.get("data")[0].get("token")
- print("rg_api_token: 已有token!")
- if token:
- db.add(ApiTokenModel(id=rg_api_token, app_id=rg_api_token, type="platform", token=token))
- db.commit()
- print("rg_api_token: 更新成功!")
- except Exception as e:
- print(e)
- db.rollback()
- async def sync_complex_api_token(db):
- token = ""
- try:
- complex_list = db.query(ComplexChatModel).all()
- for i in complex_list:
- user_token = db.query(ApiTokenModel).filter(ApiTokenModel.app_id == i.id).first()
- if not user_token:
- try:
- chat = ChatBaseApply()
- url = settings.dify_base_url + DF_CHAT_API_KEY.format(i.id)
- access_token = await get_app_token(db, workflow_server)
- param = await chat.chat_get(url, {}, await chat.get_headers(access_token))
- if param and param.get("data"):
- token = param.get("data", [{}])[0].get("token")
- token_id = param.get("data", [{}])[0].get("id")
- # dialog.parameters = json.dumps(param)
- else:
- param = await chat.chat_post(url, {}, await chat.get_headers(access_token))
- if param:
- token = param.get("token")
- token_id = param.get("id")
- if token:
- db.add(ApiTokenModel(id=token_id, app_id=i.id, type="app", token=token))
- db.commit()
- print("df_api_token: 更新成功!")
- except Exception as e:
- print(e)
- except Exception as e:
- print(e)
- db.rollback()
- async def system_license_sync(db):
- with open(os.path.join(ENV_CONF_PATH, "system.yaml") , 'r', encoding='utf-8') as file:
- # 加载JSON数据
- config = yaml.safe_load(file)
- try:
- system = db.query(SystemDataModel).filter_by(id=SYSTEM_ID).first()
- if system:
- system.version = config["smart_system"].get("version")
- else:
- system = SystemDataModel(id=SYSTEM_ID, version=config["smart_system"].get("version"), title=config["smart_system"].get("title"), desc=config["smart_system"].get("desc"), machine_id=get_machine_id())
- db.add(system)
- db.commit()
- except Exception as e:
- print(e)
- db.rollback()
|