From 799c2381e4906d103b94049df61876fc4ab7b692 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E9=9D=92=E8=BD=A9?= Date: Mon, 12 Aug 2024 16:10:47 +0800 Subject: [PATCH] log work flow --- memoryscope/core/operation/base_workflow.py | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/memoryscope/core/operation/base_workflow.py b/memoryscope/core/operation/base_workflow.py index ea11d8d9..de3c41fc 100644 --- a/memoryscope/core/operation/base_workflow.py +++ b/memoryscope/core/operation/base_workflow.py @@ -31,7 +31,7 @@ class BaseWorkflow(object): self.context: Dict[str, Any] = {} self.context_lock = threading.Lock() - self.logger: Logger = Logger.get_logger() + self.logger: Logger = Logger.get_logger(Logger.append_timestamp("workflow")) if self.workflow: self.workflow_worker_list = self._parse_workflow() @@ -165,21 +165,27 @@ class BaseWorkflow(object): **kwargs: Additional keyword arguments to be passed to context. """ with Timer(f"workflow.{self.name}", time_log_type="wrap"): + self.logger.info(f"\n\n\n++++++++++++++++++++++++ [{self.name}] ++++++++++++++++++++++++") + self.context.clear() - self.context.update({WORKFLOW_NAME: self.name, **kwargs}) - + n_stage = len(self.workflow_worker_list) # Iterate over each part of the workflow - for workflow_part in self.workflow_worker_list: + for index, workflow_part in enumerate(self.workflow_worker_list): # Sequential execution for single-item parts if len(workflow_part) == 1: + self.logger.info(f"\n-----------------------------------------------------") + self.logger.info(f"sequential execution ({self.name}) | {index+1}/{n_stage}: {workflow_part[0]}") if not self._run_sub_workflow(workflow_part[0]): break # Parallel execution for multi-item parts else: t_list = [] # Submit tasks to the thread pool - for sub_workflow in workflow_part: + n_sub_stage = len(workflow_part) + for sub_index, sub_workflow in enumerate(workflow_part): + self.logger.info(f"\n-----------------------------------------------------") + self.logger.info(f"parallel sequential execution ({self.name}) | {index+1}/{n_stage} | {sub_index+1}/{n_sub_stage}: {str(sub_workflow)}") t_list.append(self.thread_pool.submit(self._run_sub_workflow, sub_workflow)) # Check results; if any task returns False, stop the workflow