mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-07-25 09:53:29 +08:00
@@ -55,16 +55,11 @@ class PipelineOperationLogService(CommonService):
|
||||
model = PipelineOperationLog
|
||||
|
||||
@classmethod
|
||||
def _final_operation_statuses(cls, operation_status):
|
||||
final_statuses = [TaskStatus.DONE.value, TaskStatus.FAIL.value]
|
||||
if not operation_status:
|
||||
return final_statuses
|
||||
requested = {status.value if isinstance(status, TaskStatus) else str(status) for status in operation_status}
|
||||
return [status for status in final_statuses if status in requested]
|
||||
|
||||
@classmethod
|
||||
def _is_final_progress(cls, progress):
|
||||
return progress == 1 or progress == -1
|
||||
def _is_final_state(cls, progress, operation_status):
|
||||
if progress == 1 or progress == -1:
|
||||
return True
|
||||
status = operation_status.value if isinstance(operation_status, TaskStatus) else str(operation_status)
|
||||
return status in [TaskStatus.CANCEL.value, TaskStatus.DONE.value, TaskStatus.FAIL.value]
|
||||
|
||||
@classmethod
|
||||
def get_file_logs_fields(cls):
|
||||
@@ -175,13 +170,13 @@ class PipelineOperationLogService(CommonService):
|
||||
raise RuntimeError(f"Task not found for dataset {document.kb_id}")
|
||||
title = task_type
|
||||
document_name = task_type
|
||||
operation_status = TaskStatus.DONE if task.progress == 1 else TaskStatus.FAIL
|
||||
operation_status = TaskStatus.DONE.value if task.progress == 1 else TaskStatus.FAIL.value if task.progress == -1 else TaskStatus.RUNNING.value
|
||||
progress = task.progress
|
||||
progress_msg = task.progress_msg
|
||||
process_begin_at = task.begin_at
|
||||
process_duration = task.process_duration
|
||||
|
||||
if not cls._is_final_progress(progress):
|
||||
if not cls._is_final_state(progress, operation_status):
|
||||
logging.info("Skip non-final dataset pipeline operation log task_id=%s task_type=%s progress=%s", task_id, task_type, progress)
|
||||
return None
|
||||
|
||||
@@ -190,7 +185,7 @@ class PipelineOperationLogService(CommonService):
|
||||
document.kb_id,
|
||||
{_PIPELINE_TASK_TYPE_TO_FINISH_FIELD[task_type]: finish_at},
|
||||
)
|
||||
elif not cls._is_final_progress(progress):
|
||||
elif not cls._is_final_state(progress, operation_status):
|
||||
logging.info("Skip non-final file pipeline operation log document_id=%s task_type=%s progress=%s", document_id, task_type, progress)
|
||||
return None
|
||||
|
||||
@@ -251,7 +246,8 @@ class PipelineOperationLogService(CommonService):
|
||||
|
||||
logs = logs.where(cls.model.document_id != GRAPH_RAPTOR_FAKE_DOC_ID)
|
||||
|
||||
logs = logs.where(cls.model.operation_status.in_(cls._final_operation_statuses(operation_status)))
|
||||
if operation_status:
|
||||
logs = logs.where(cls.model.operation_status.in_([status.value if isinstance(status, TaskStatus) else str(status) for status in operation_status]))
|
||||
if types:
|
||||
logs = logs.where(cls.model.document_type.in_(types))
|
||||
if suffix:
|
||||
@@ -287,7 +283,8 @@ class PipelineOperationLogService(CommonService):
|
||||
else:
|
||||
logs = cls.model.select(*fields).where((cls.model.kb_id == kb_id), (cls.model.document_id == GRAPH_RAPTOR_FAKE_DOC_ID))
|
||||
|
||||
logs = logs.where(cls.model.operation_status.in_(cls._final_operation_statuses(operation_status)))
|
||||
if operation_status:
|
||||
logs = logs.where(cls.model.operation_status.in_([status.value if isinstance(status, TaskStatus) else str(status) for status in operation_status]))
|
||||
if create_date_from:
|
||||
logs = logs.where(cls.model.create_date >= create_date_from)
|
||||
if create_date_to:
|
||||
|
||||
Reference in New Issue
Block a user