feat: workflow chat support continue exec

This commit is contained in:
GuoQing Zhang
2024-12-24 15:53:49 +08:00
parent 25aad25158
commit fea447c8d1
4 changed files with 20 additions and 17 deletions
+1 -9
View File
@@ -62,19 +62,11 @@ logger_conf:
# 日志级别
level: INFO
# 日志格式化函数,extra内支持trace_id
format: "[{time:YYYY-MM-DD HH:mm:ss.SSSSSS}]|{level}|BISHENG|{extra[trace_id]}|{process.id}|{thread.id}|{message}"
format: '<level>[{time:YYYY-MM-DD HH:mm:ss.SSSSSS}] [{level.name} process-{process.id}-{thread.id} {name}:{line}]</level> - <level>trace={extra[trace_id]} {message}</level>'
# 每天的几点进行切割
rotation: "00:00"
retention: "3 Days"
enqueue: ture
- sink: "/app/data/err-v0-BISHENG-{HOSTNAME}.log"
level: ERROR
# 和原生不一样,后端会将配置使用eval()执行转为函数用来过滤特定日志级别。推荐lambda
filter: "lambda record: record['level'].name == 'ERROR'"
format: "[{time:YYYY-MM-DD HH:mm:ss.SSSSSS}]|{level}|BISHENG|{extra[trace_id]}||{process.id}|{thread.id}|||#EX_ERR:POS={name},line {line},ERR=500,EMSG={message}"
rotation: "00:00"
retention: "3 Days"
enqueue: ture
- sink: "/app/data/statistic.log"
level: INFO
# 和原生不一样,后端会将配置使用eval()执行转为函数用来过滤特定日志级别。推荐lambda
+4 -1
View File
@@ -48,7 +48,10 @@ class RedisClient:
try:
if pickled := pickle.dumps(value):
self.cluster_nodes(key)
result = self.connection.setex(key, expiration, pickled)
if expiration:
result = self.connection.setex(key, expiration, pickled)
else:
result = self.connection.set(key, pickled)
if not result:
raise ValueError('RedisCache could not set the value.')
else:
@@ -27,11 +27,13 @@ class WorkflowClient(BaseClient):
websocket, **kwargs)
self.workflow: Optional[RedisCallback] = None
self.history = []
async def close(self):
if self.workflow:
# 非会话模式关闭workflow执行
if self.workflow and not self.chat_id:
self.workflow.set_workflow_stop()
self.workflow = None
self.workflow = None
async def save_chat_message(self, chat_response: ChatResponse) -> int | None:
if not self.chat_id:
@@ -75,8 +77,8 @@ class WorkflowClient(BaseClient):
async def init_history(self):
if not self.chat_id:
return
history = ChatMessageDao.get_latest_message_by_chatid(self.chat_id)
if not history:
self.history = ChatMessageDao.get_latest_message_by_chatid(self.chat_id)
if not self.history:
# 新建会话,记录审计日志
AuditLogService.create_chat_workflow(self.login_user, get_request_ip(self.request), self.client_id)
@@ -95,10 +97,13 @@ class WorkflowClient(BaseClient):
self.workflow = RedisCallback(unique_id, workflow_id, self.chat_id, str(self.user_id))
status_info = self.workflow.get_workflow_status()
if not status_info:
if self.history:
await self.send_response('processing', 'close', '')
self.workflow.set_workflow_data(workflow_data)
self.workflow.set_workflow_status(WorkflowStatus.WAITING.value)
# 发起异步任务
execute_workflow.delay(unique_id, workflow_id, self.chat_id, str(self.user_id))
except Exception as e:
logger.exception('init_workflow_error')
self.workflow = None
@@ -136,10 +141,11 @@ class WorkflowClient(BaseClient):
else:
send_msg = False
self.workflow = None
if status_info['status'] == WorkflowStatus.FAILED.value:
await self.send_response('error', 'over', status_info['reason'])
await self.send_response('processing', 'close', '')
self.workflow.clear_workflow_status()
self.workflow = None
break
else:
chat_response = self.workflow.get_workflow_response()
@@ -46,8 +46,7 @@ class RedisCallback(BaseCallback):
def set_workflow_status(self, status: int, reason: str = None):
self.redis_client.set(self.workflow_status_key,
{'status': status, 'reason': reason, 'time': time.time()},
expiration=self.workflow_expire_time)
self.workflow_cache.clear()
expiration=None)
if status in [WorkflowStatus.FAILED.value, WorkflowStatus.SUCCESS.value]:
# 消息事件和状态key可能还需要消费
self.redis_client.delete(self.workflow_data_key)
@@ -60,6 +59,9 @@ class RedisCallback(BaseCallback):
self.workflow_cache.setdefault(self.workflow_status_key, workflow_status)
return workflow_status
def clear_workflow_status(self):
self.redis_client.delete(self.workflow_status_key)
def insert_workflow_response(self, event: dict):
self.redis_client.rpush(self.workflow_event_key, json.dumps(event), expiration=self.workflow_expire_time)