initialize_data.py 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292
  1. import json
  2. import time
  3. import os
  4. import yaml
  5. from passlib.context import CryptContext
  6. from Log import logger
  7. # from app.api import pwd_context
  8. from app.config.const import DIFY, ENV_CONF_PATH, RAGFLOW
  9. from app.models import MenuCapacityModel, WebMenuModel, GroupModel, RoleModel, DialogModel, UserModel, UserAppModel, \
  10. cipher_suite, UserTokenModel
  11. from app.service.auth import UserAppDao
  12. from app.service.bisheng import BishengService
  13. from app.service.difyService import DifyService
  14. from app.service.ragflow import RagflowService
  15. from app.service.service_token import get_new_token
  16. from app.service.v2.app_register import AppRegisterDao
  17. from app.config.config import settings
  18. from app.utils.password_handle import generate_password
  19. pwd_context = CryptContext(schemes=["bcrypt"], deprecated="auto")
  20. async def dialog_menu_sync(db):
  21. menu_list = []
  22. with open(os.path.join(ENV_CONF_PATH, "menu_conf.json") , 'r', encoding='utf-8') as file:
  23. # 加载JSON数据
  24. data = json.load(file)
  25. menu_list = data.get("data", [])
  26. db.query(WebMenuModel).delete()
  27. db.query(MenuCapacityModel).delete()
  28. db.commit()
  29. for menu in menu_list:
  30. # print(menu)
  31. dialog = menu.pop("dialog", [])
  32. for i in dialog:
  33. capacity = MenuCapacityModel(menu_id=menu["id"], capacity_id=i["id"], capacity_type=i["agentType"],
  34. chat_id=i["id"] if not i["chat_id"] else i["chat_id"],
  35. chat_type=i["chat_type"])
  36. db.add(capacity)
  37. menu_obj = WebMenuModel(**menu)
  38. db.add(menu_obj)
  39. db.commit()
  40. async def create_menu_sync(db):
  41. # json_file_path = "env_conf/menu_conf.json.template"
  42. json_file_path = os.path.join(ENV_CONF_PATH, "menu_conf.json.template")
  43. with open(json_file_path, 'r', encoding='utf-8') as file:
  44. json_data = json.load(file).get("data", [])
  45. # for menu in json_data:
  46. # menu['dialog'].clear()
  47. dialogs = db.query(DialogModel).all()
  48. dialog_dict = {}
  49. for dialog in dialogs:
  50. if dialog.name not in dialog_dict:
  51. dialog_dict[dialog.name] = []
  52. dialog_dict[dialog.name].append({
  53. 'id': dialog.id,
  54. 'chat_id': dialog.id,
  55. 'chat_type': '',
  56. 'agentType': dialog.dialog_type
  57. })
  58. for menu in json_data:
  59. # if menu['title'] in dialog_dict:
  60. # for dialog in dialog_dict[menu['title']]:
  61. # new_dialog_item = {
  62. # 'id': dialog.id,
  63. # 'chat_id': dialog.id,
  64. # 'chat_type': '',
  65. # 'agentType': dialog.dialog_type
  66. # }
  67. menu['dialog']= dialog_dict.get(menu['title'], [])
  68. json_data = {"data": json_data}
  69. new_file_name = f"menu_conf.json.template"
  70. new_file_path = os.path.join(os.path.dirname(json_file_path), new_file_name)
  71. with open(new_file_path, 'w', encoding='utf-8') as new_file:
  72. json.dump(json_data, new_file, ensure_ascii=False, indent=4)
  73. return {
  74. "file_name": new_file_name,
  75. "json_data": json_data
  76. }
  77. async def default_group_sync(db):
  78. group = db.query(GroupModel).filter_by(group_type=2).first()
  79. if not group:
  80. logger.error("未初始默认组, 开始初始化!")
  81. try:
  82. group = GroupModel(group_name="默认用户组", group_description="默认组", group_type=2)
  83. db.add(group)
  84. db.commit()
  85. except Exception as e:
  86. logger.error(e)
  87. async def default_role_sync(db):
  88. role = db.query(RoleModel).filter_by(role_type=2).first()
  89. if not role:
  90. logger.error("未初始默认角色, 开始初始化!")
  91. try:
  92. group = RoleModel(id="morenjuese1234567890", name="默认角色", description="默认角色", role_type=2)
  93. db.add(group)
  94. db.commit()
  95. except Exception as e:
  96. logger.error(e)
  97. async def app_register_sync(db):
  98. app_dict = {}
  99. with open(os.path.join(ENV_CONF_PATH, "app_register_conf.json"), 'r', encoding='utf-8') as file:
  100. # 加载JSON数据
  101. app_dict = json.load(file)
  102. try:
  103. for app_id, status in app_dict.items():
  104. AppRegisterDao(db).update_and_insert_app(app_id, status)
  105. except Exception as e:
  106. logger.error(e)
  107. async def basic_agent_sync(db):
  108. agent_list = []
  109. with open(os.path.join(ENV_CONF_PATH, "default_agent_conf.json"), 'r', encoding='utf-8') as file:
  110. # 加载JSON数据
  111. agent_dict = json.load(file)
  112. agent_list = agent_dict.get("basic", [])
  113. user = db.query(UserModel).filter_by(permission="admin").first()
  114. for agent in agent_list:
  115. dialog = db.query(DialogModel).filter(DialogModel.id == agent["id"]).first()
  116. if dialog:
  117. try:
  118. dialog.name = agent["name"]
  119. dialog.description = agent["description"]
  120. dialog.icon = agent["icon"]
  121. dialog.mode = agent["mode"]
  122. dialog.parameters = json.dumps(agent["parameters"])
  123. db.commit()
  124. except Exception as e:
  125. logger.error(e)
  126. else:
  127. try:
  128. dialog = DialogModel(id=agent["id"], name=agent["name"], description=agent["description"],
  129. icon=agent["icon"], tenant_id=user.id if user else "", dialog_type=agent["dialogType"], mode=agent["mode"],parameters = json.dumps(agent["parameters"]))
  130. db.add(dialog)
  131. db.commit()
  132. db.refresh(dialog)
  133. except Exception as e:
  134. print(e)
  135. db.rollback()
  136. async def user_update_app(userid, db):
  137. user = db.query(UserModel).filter(UserModel.id == userid).first()
  138. if not user:
  139. raise Exception("User id not found")
  140. app_register = AppRegisterDao(db).get_apps()
  141. register_dict = {}
  142. token = ""
  143. app_password = await generate_password(10)
  144. crypt_password = UserAppModel.encrypted_password(app_password)
  145. for app in app_register:
  146. if app["id"] == 'ragflow_app':
  147. user_rag_app = db.query(UserAppModel).filter(UserAppModel.user_id == userid,
  148. UserAppModel.app_type == 'ragflow_app').all()
  149. if not user_rag_app:
  150. service = RagflowService(settings.fwr_base_url)
  151. register_info = await register_app(service, app["id"], app_password, token)
  152. if register_info:
  153. register_dict[app["id"]] = register_info
  154. app_name = register_info.get("name")
  155. app_id = register_info.get("id")
  156. app_email = register_info.get("email")
  157. await save_db(db, app_name, crypt_password, app_email, user.id, app_id, "ragflow_app")
  158. elif app["id"] == 'bisheng_app':
  159. user_bs_app = db.query(UserAppModel).filter(UserAppModel.user_id == userid,
  160. UserAppModel.app_type == 'bisheng_app').all()
  161. if not user_bs_app:
  162. service = BishengService(settings.sgb_base_url)
  163. register_info = await register_app(service, app["id"], app_password, token)
  164. if register_info:
  165. register_dict[app["id"]] = register_info
  166. app_name = register_info.get("name")
  167. app_id = register_info.get("id")
  168. app_email = register_info.get("email")
  169. await save_db(db, app_name, crypt_password, app_email, user.id, app_id, "bisheng_app")
  170. elif app["id"] == 'dify_app':
  171. user_df_app = db.query(UserAppModel).filter(UserAppModel.user_id == userid,
  172. UserAppModel.app_type == 'dify_app').all()
  173. if not user_df_app:
  174. admin_user = db.query(UserModel).filter(UserModel.permission == "admin").first()
  175. token = await get_new_token(db, admin_user.id, DIFY)
  176. if not token:
  177. print("用户注册获取dftoken失败!")
  178. service = DifyService(settings.dify_base_url)
  179. register_info = await register_app(service, app["id"], app_password, token)
  180. if register_info:
  181. register_dict[app["id"]] = register_info
  182. app_name = register_info.get("name")
  183. app_id = register_info.get("id")
  184. app_email = register_info.get("email")
  185. await save_db(db, app_name, crypt_password, app_email, user.id, app_id, "dify_app")
  186. else:
  187. raise Exception("未知注册应用---")
  188. async def register_app(service, app_id, app_password, token):
  189. name = app_id + str(int(time.time()))
  190. try:
  191. register_info = await service.register(name, app_password, token)
  192. return {"id": register_info.get("id"), "name": name, "email": register_info.get("email")}
  193. except Exception as e:
  194. print(f"Failed to register with {app_id}: {str(e)}")
  195. return None
  196. async def save_db(db, username, password, email, user_id, app_id, app_type):
  197. user_app_dao = UserAppDao(db)
  198. user_id = await user_app_dao.insert_user_app_data(username, password, email, user_id, app_id, app_type)
  199. if not user_id:
  200. raise Exception("Failed to register with app")
  201. print({"msg": "User registered successfully", "userFlag": user_id})
  202. async def admin_account_sync(db):
  203. try:
  204. config = {}
  205. now_account =[]
  206. with open(os.path.join(ENV_CONF_PATH, "account.yaml"), 'r', encoding='utf-8') as file:
  207. # 加载JSON数据
  208. config = yaml.safe_load(file)
  209. account_list = db.query(UserTokenModel).all()
  210. for account in account_list:
  211. if account.id in config:
  212. if account.account != config[account.id]["account"] or account.password != config[account.id]["password"]:
  213. db.query(UserTokenModel).filter_by(id=account.id).update({"account": config[account.id]["account"],
  214. "password": config[account.id]["password"],
  215. "access_token": ""
  216. })
  217. now_account.append(account.id)
  218. else:
  219. db.query(UserTokenModel).filter_by(id=account.id).delete()
  220. for k, v in config.items():
  221. if k not in now_account:
  222. new_account = UserTokenModel(id=k, account=v["account"], password=v["password"])
  223. db.add(new_account)
  224. db.commit()
  225. except Exception as e:
  226. print(e)
  227. db.rollback()
  228. async def admin_user_sync(db):
  229. try:
  230. config = {}
  231. with open(os.path.join(ENV_CONF_PATH, "admin.yaml"), 'r', encoding='utf-8') as file:
  232. # 加载JSON数据
  233. config = yaml.safe_load(file)
  234. # print(config)
  235. db_user = db.query(UserModel).filter(UserModel.username == config["smart_server"]["account"]).first()
  236. if db_user:
  237. print("admin_user_sync: 用户已经存在!")
  238. return
  239. register_dict = {}
  240. for app in [RAGFLOW, DIFY]:
  241. register_dict[app] = {"id": config[app].get("id", "123"), "name": config[app]["account"],
  242. "pwd":config[app]["password"],
  243. "email": config[app]["account"]}
  244. # 存储用户信息
  245. hashed_password = pwd_context.hash(config["smart_server"]["password"])
  246. user_model = UserModel(username=config["smart_server"]["account"], hashed_password=hashed_password, email="",
  247. phone="", login_name="", sync_flag="", creator=0, permission="admin")
  248. db.add(user_model)
  249. db.commit()
  250. db.refresh(user_model)
  251. u_id = user_model.id
  252. user_app_dao = UserAppDao(db)
  253. for k, v in register_dict.items():
  254. 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)
  255. except Exception as e:
  256. print(e)
  257. db.rollback()