| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256 |
- import json
- from datetime import datetime
- from sqlalchemy import or_
- from app.config.agent_base_url import DF_CHAT_PARAMETERS, DF_CHAT_API_KEY
- from app.config.config import settings
- from app.config.const import Dialog_STATSU_DELETE, DF_TYPE, Dialog_STATSU_ON, workflow_server, RG_TYPE
- from app.models import KnowledgeModel, GroupModel, DialogModel, ConversationModel, group_dialog_table, LabelWorkerModel, \
- LabelModel, ApiTokenModel
- from app.models.user_model import UserModel, UserTokenModel
- from Log import logger
- from app.service.v2.app_driver.chat_data import ChatBaseApply
- from app.service.v2.chat import get_chat_token, add_chat_token, get_app_token
- from app.task.fetch_agent import get_one_from_ragflow_dialog
- async def get_dialog_list(db, user_id, keyword, label, status, page_size, page_index):
- user = db.query(UserModel).filter(UserModel.id == user_id).first()
- if user is None:
- return {"rows": []}
- query = db.query(DialogModel)
- if status:
- query = query.filter(DialogModel.status == status)
- else:
- query = query.filter(DialogModel.status != Dialog_STATSU_DELETE)
- id_list = []
- # if label:
- # id_list = [i.object_id for i in db.query(LabelWorkerModel).filter(LabelWorkerModel.label_id==label).all()]
- if user.permission != "admin":
- dia_list = [j.id for i in user.groups for j in i.dialogs]
- query = query.filter(or_(DialogModel.tenant_id == user_id, DialogModel.id.in_(dia_list)))
- # else:
- if label:
- id_list = [i.object_id for i in db.query(LabelWorkerModel).filter(LabelWorkerModel.label_id == label).all()]
- query = query.filter(DialogModel.id.in_(id_list))
- if keyword:
- query = query.filter(DialogModel.name.like('%{}%'.format(keyword)))
- query = query.order_by(DialogModel.update_date.desc())
- total = query.count()
- if page_size:
- query = query.limit(page_size).offset((page_index - 1) * page_size)
- rows = []
- user_id_set = set()
- dialog_id_set = set()
- label_dict = {}
- for kld in query.all():
- user_id_set.add(kld.tenant_id)
- dialog_id_set.add(kld.id)
- rows.append(kld.to_json())
- user_dict = {str(i.id): i.to_dict() for i in db.query(UserModel).filter(UserModel.id.in_(user_id_set)).all()}
- for i in db.query(LabelModel.id, LabelModel.name, LabelWorkerModel.object_id).outerjoin(LabelWorkerModel,
- LabelModel.id == LabelWorkerModel.label_id).filter(
- LabelWorkerModel.object_id.in_(dialog_id_set)).all():
- label_dict[i.object_id] = label_dict.get(i.object_id, []) + [{"labelId": i.id, "labelName": i.name}]
- for r in rows:
- r["user"] = user_dict.get(r["user_id"], {})
- r["label"] = label_dict.get(r["id"], [])
- return {"total": total, "rows": rows}
- async def update_session_history(db, data: dict, user_id):
- session_id = data.get("id")
- if not session_id:
- logger.error("更新回话记录失败!{}".format(data))
- return
- data["create_date"] = datetime.strptime(data["create_date"], '%a, %d %b %Y %H:%M:%S %Z')
- data["update_date"] = datetime.strptime(data["update_date"], '%a, %d %b %Y %H:%M:%S %Z')
- conversation = db.query(ConversationModel).filter(ConversationModel.id == session_id).first()
- if not conversation:
- try:
- data["tenant_id"] = user_id
- conversation_model = ConversationModel(**data)
- db.add(conversation_model)
- db.commit()
- except Exception as e:
- logger.error(e)
- db.rollback()
- else:
- try:
- # data["tenant_id"] = user_id
- del data["id"]
- db.query(ConversationModel).filter(ConversationModel.id == session_id).update(data)
- db.commit()
- except Exception as e:
- logger.error(e)
- db.rollback()
- async def get_session_history(db, user_id, dialog_id, page, limit):
- session_list = db.query(ConversationModel).filter(ConversationModel.tenant_id.__eq__(user_id),
- ConversationModel.dialog_id.__eq__(dialog_id)).order_by(
- ConversationModel.update_time.desc()).limit(limit).offset((page - 1) * limit).all()
- return [i.to_json() for i in session_list]
- async def create_dialog_service(db, dialog_id, dialog_name, description, icon, dialog_type, mode, user_id):
- para = {
- "user_input_form": [],
- "retriever_resource": {
- "enabled": True
- },
- "file_upload": {
- "enabled": False
- }
- }
- try:
- dialog_model = DialogModel(id=dialog_id, name=dialog_name, description=description, icon=icon,
- dialog_type=dialog_type, tenant_id=user_id, mode=mode, update_date=datetime.now(),
- create_date=datetime.now(), parameters=json.dumps(para))
- db.add(dialog_model)
- db.commit()
- db.refresh(dialog_model)
- except Exception as e:
- logger.error(e)
- db.rollback()
- return False
- return True
- async def update_dialog_status_service(db, dialog_id, status, user_id):
- try:
- dialog = db.query(DialogModel).filter_by(id=dialog_id).first()
- dialog.status = status
- dialog.update_date = datetime.now()
- # db.query(DialogModel).filter_by(id=dialog_id).update({"status":status, "update_date": datetime.now()})
- if dialog.dialog_type == DF_TYPE and status == Dialog_STATSU_ON:
- chat = ChatBaseApply()
- token = await get_chat_token(db, dialog_id)
- if not token:
- access_token = await get_app_token(db, workflow_server)
- # print(workflow)
- if access_token:
- url = settings.dify_base_url + DF_CHAT_API_KEY.format(dialog_id)
- 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")
- await add_chat_token(db, {"id":token_id, "app_id": dialog_id, "type":"app", "token": token})
- # 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")
- await add_chat_token(db, {"id": token_id, "app_id": dialog_id, "type": "app", "token": token})
- if token:
- url = settings.dify_base_url + DF_CHAT_PARAMETERS
- param = await chat.chat_get(url, {"user": str(user_id)}, await chat.get_headers(token))
- if param:
- dialog.parameters = json.dumps(param)
- db.commit()
- except Exception as e:
- logger.error(e)
- db.rollback()
- return False
- return True
- async def delete_dialog_service(db, dialog_id):
- try:
- db.query(DialogModel).filter_by(id=dialog_id).update(
- {"status": Dialog_STATSU_DELETE, "update_date": datetime.now()})
- db.commit()
- except Exception as e:
- logger.error(e)
- db.rollback()
- return False
- return True
- async def update_dialog_icon_service(db, dialog_id, icon, name, description):
- update = {"icon": icon, "update_date": datetime.now()}
- if name:
- update["name"] = name
- if description or description == "":
- update["description"] = description
- try:
- db.query(DialogModel).filter_by(id=dialog_id).update(update)
- db.commit()
- except Exception as e:
- logger.error(e)
- db.rollback()
- return False
- return True
- async def get_dialog_manage_list(db, user_id, keyword, label, status, page_size, page_index, mode):
- user = db.query(UserModel).filter(UserModel.id == user_id).first()
- if user is None:
- return {"rows": []}
- query = db.query(DialogModel).filter(DialogModel.status != Dialog_STATSU_DELETE)
- if user.permission != "admin":
- dia_list = [j.id for i in user.groups for j in i.dialogs]
- query = query.filter(or_(DialogModel.tenant_id == user_id, DialogModel.id.in_(dia_list)))
- if label:
- id_list = set(
- [i.object_id for i in db.query(LabelWorkerModel).filter(LabelWorkerModel.label_id.in_(label)).all()])
- query = query.filter(DialogModel.id.in_(id_list))
- if keyword:
- query = query.filter(DialogModel.name.like('%{}%'.format(keyword)))
- if status:
- # print(status)
- query = query.filter(DialogModel.status == status)
- if mode:
- query = query.filter(DialogModel.mode == mode)
- query = query.order_by(DialogModel.update_date.desc())
- total = query.count()
- if page_size:
- query = query.limit(page_size).offset((page_index - 1) * page_size)
- rows = []
- user_id_set = set()
- dialog_id_set = set()
- label_dict = {}
- for kld in query.all():
- user_id_set.add(kld.tenant_id)
- dialog_id_set.add(kld.id)
- rows.append(kld.to_json())
- user_dict = {str(i.id): i.to_dict() for i in db.query(UserModel).filter(UserModel.id.in_(user_id_set)).all()}
- for i in db.query(LabelModel.id, LabelModel.name, LabelWorkerModel.object_id).outerjoin(LabelWorkerModel,
- LabelModel.id == LabelWorkerModel.label_id).filter(
- LabelWorkerModel.object_id.in_(dialog_id_set)).all():
- label_dict[i.object_id] = label_dict.get(i.object_id, []) + [{"labelId": i.id, "labelName": i.name}]
- for r in rows:
- r["user"] = user_dict.get(r["user_id"], {})
- r["label"] = label_dict.get(r["id"], [])
- return {"total": total, "rows": rows}
- async def sync_dialog_service(db, dialog_id):
- dialog = db.query(DialogModel).filter(DialogModel.id == dialog_id).first()
- if dialog and dialog.dialog_type == RG_TYPE:
- try:
- app_dialog = get_one_from_ragflow_dialog(dialog_id)
- if app_dialog:
- dialog.name = app_dialog["name"]
- dialog.description = app_dialog["description"]
- dialog.update_date = datetime.now()
- db.add(dialog)
- db.commit()
- db.refresh(dialog)
- except Exception as e:
- logger.error(e)
- db.rollback()
- return False
- return True
|