mirror of
https://github.com/infiniflow/ragflow.git
synced 2026-07-25 01:43:27 +08:00
@@ -54,6 +54,18 @@ _PIPELINE_TASK_TYPE_TO_FINISH_FIELD = {
|
||||
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
|
||||
|
||||
@classmethod
|
||||
def get_file_logs_fields(cls):
|
||||
return [
|
||||
@@ -169,11 +181,18 @@ class PipelineOperationLogService(CommonService):
|
||||
process_begin_at = task.begin_at
|
||||
process_duration = task.process_duration
|
||||
|
||||
if not cls._is_final_progress(progress):
|
||||
logging.info("Skip non-final dataset pipeline operation log task_id=%s task_type=%s progress=%s", task_id, task_type, progress)
|
||||
return None
|
||||
|
||||
finish_at = process_begin_at + timedelta(seconds=process_duration)
|
||||
KnowledgebaseService.update_by_id(
|
||||
document.kb_id,
|
||||
{_PIPELINE_TASK_TYPE_TO_FINISH_FIELD[task_type]: finish_at},
|
||||
)
|
||||
elif not cls._is_final_progress(progress):
|
||||
logging.info("Skip non-final file pipeline operation log document_id=%s task_type=%s progress=%s", document_id, task_type, progress)
|
||||
return None
|
||||
|
||||
log = dict(
|
||||
id=get_uuid(),
|
||||
@@ -232,8 +251,7 @@ class PipelineOperationLogService(CommonService):
|
||||
|
||||
logs = logs.where(cls.model.document_id != GRAPH_RAPTOR_FAKE_DOC_ID)
|
||||
|
||||
if operation_status:
|
||||
logs = logs.where(cls.model.operation_status.in_(operation_status))
|
||||
logs = logs.where(cls.model.operation_status.in_(cls._final_operation_statuses(operation_status)))
|
||||
if types:
|
||||
logs = logs.where(cls.model.document_type.in_(types))
|
||||
if suffix:
|
||||
@@ -269,8 +287,7 @@ 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))
|
||||
|
||||
if operation_status:
|
||||
logs = logs.where(cls.model.operation_status.in_(operation_status))
|
||||
logs = logs.where(cls.model.operation_status.in_(cls._final_operation_statuses(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