Loading csst_common/pipeline.py +93 −36 Original line number Diff line number Diff line Loading @@ -60,8 +60,9 @@ class Pipeline: """ Examples -------- >>> p = Pipeline(msg="{'a':1}") >>> p.msg_dict >>> p = Pipeline() >>> p.info() >>> p.task {'a': 1} """ Loading Loading @@ -123,6 +124,18 @@ class Pipeline: self.task = {} self.logger.info(f"task = {self.task}") # warning operations @staticmethod def filter_warnings(level: str = "ignore"): # Suppress all warnings warnings.filterwarnings(level) @staticmethod def reset_warnings(self): """Reset warning filters.""" warnings.resetwarnings() # message operations @staticmethod def dict2json(d: dict): """Convert `dict` to JSON format string.""" Loading @@ -133,18 +146,12 @@ class Pipeline: """Convert JSON format string to `dict`.""" return json.loads(m) @property def summarize(self): """Summarize this run.""" t_stop: Time = Time.now() t_cost: float = (t_stop - self.t_start).value * 86400.0 self.logger.info(f"Total cost: {t_cost:.1f} sec") def clean_output(self): """Clean output directory.""" # self.clean_directory(self.dir_output) # clean in run command pass # file operations @staticmethod def mkdir(d): """Create a directory if it does not exist.""" if not os.path.exists(d): os.makedirs(d) @staticmethod def clean_directory(d): Loading @@ -158,19 +165,26 @@ class Pipeline: print("Failed to clean output directory!") print_directory_tree(d) def now(self): """Return ISOT format datetime using `astropy`.""" return Time.now().isot @staticmethod def filter_warnings(level: str = "ignore"): # Suppress all warnings warnings.filterwarnings(level) def move(file_src: str, file_dst: str) -> str: """Move file `file_src` to `file_dist`.""" return shutil.move(file_src, file_dst) @staticmethod def reset_warnings(self): """Reset warning filters.""" warnings.resetwarnings() def copy(file_src: str, file_dst: str) -> str: """Move file `file_src` to `file_dist`.""" return shutil.copy(file_src, file_dst) @property def summarize(self): """Summarize this run.""" t_stop: Time = Time.now() t_cost: float = (t_stop - self.t_start).value * 86400.0 self.logger.info(f"Total cost: {t_cost:.1f} sec") def clean_output(self): """Clean output directory.""" self.clean_directory(self.dir_output) def file(self, file_path): """Initialize File object.""" Loading @@ -180,25 +194,66 @@ class Pipeline: """Create new file in output directory.""" return os.path.join(self.dir_output, file_name) @staticmethod def move(file_src: str, file_dst: str): """Move file `file_src` to `file_dist`.""" shutil.move(file_src, file_dst) @staticmethod def copy(file_src: str, file_dst: str): """Move file `file_src` to `file_dist`.""" shutil.copy(file_src, file_dst) def download_oss_file(self, oss_file_path: str) -> str: def download_oss_file(self, oss_file_path: str, dir_dst: str = None) -> str: """Download an OSS file from OSS to output directory.""" local_file_path = os.path.join(self.dir_output, os.path.basename(oss_file_path)) if dir_dst is None: dir_dst = self.dir_output local_file_path = os.path.join(dir_dst, os.path.basename(oss_file_path)) csst_fs.s3_fs.get(oss_file_path, local_file_path, s3_options=s3_options) assert os.path.exists( local_file_path ), f"Failed to download {oss_file_path} to {local_file_path}" return local_file_path def abspath(self, file_path: str) -> str: """Return absolute path of `file_path`.""" if file_path.__contains__(":"): # it's an OSS file path assert self.use_oss, "USE_OSS must be True to use OSS file path!" # download OSS file to output directory local_file_path = self.download_oss_file(file_path) # return local file path return local_file_path else: # it's a NAS file path if file_path.startswith("CSST"): # DFS return os.path.join(self.dfs_root, file_path) else: # CCDS return os.path.join(self.ccds_root, file_path) def download_dfs_file(self, file_path: str, dir_dst: str = None): """Copy DFS file to output directory.""" if dir_dst is None: dir_dst = self.dir_output if self.use_oss: # download OSS file to dst directory return self.download_oss_file(file_path, dir_dst) else: # copy DFS file to dst directory local_file_path = os.path.join(dir_dst, os.path.basename(file_path)) self.copy(self.abspath(file_path), local_file_path) return local_file_path def copy_ccds_refs(self, refs: dict, dir_dst: str = None) -> dict: """Copy raw file from CCDS to output directory.""" if dir_dst is None: dir_dst = self.dir_output local_refs = {} for ref_name, ref_path in refs.items(): local_file_path = os.path.join(dir_dst, os.path.basename(ref_path)) local_refs[ref_name] = self.copy(self.abspath(ref_path), local_file_path) return local_refs # time operations def now(self) -> str: """Return ISOT format datetime using `astropy`.""" return Time.now().isot # call modules def call(self, func: Callable, *args: Any, **kwargs: Any): self.logger.info(f"=====================================================") t_start: Time = Time.now() Loading Loading @@ -249,6 +304,7 @@ class Pipeline: output=output, ) # retry operations @staticmethod def retry(*args, **kwargs): return retry(*args, **kwargs) Loading @@ -261,6 +317,7 @@ class Pipeline: else: raise ValueError("CCDS client not initialized!") # automatically log pipeline information def info(self): """Return pipeline information.""" self.logger.info(f"DOCKER_IMAGE={self.docker_image}") Loading Loading
csst_common/pipeline.py +93 −36 Original line number Diff line number Diff line Loading @@ -60,8 +60,9 @@ class Pipeline: """ Examples -------- >>> p = Pipeline(msg="{'a':1}") >>> p.msg_dict >>> p = Pipeline() >>> p.info() >>> p.task {'a': 1} """ Loading Loading @@ -123,6 +124,18 @@ class Pipeline: self.task = {} self.logger.info(f"task = {self.task}") # warning operations @staticmethod def filter_warnings(level: str = "ignore"): # Suppress all warnings warnings.filterwarnings(level) @staticmethod def reset_warnings(self): """Reset warning filters.""" warnings.resetwarnings() # message operations @staticmethod def dict2json(d: dict): """Convert `dict` to JSON format string.""" Loading @@ -133,18 +146,12 @@ class Pipeline: """Convert JSON format string to `dict`.""" return json.loads(m) @property def summarize(self): """Summarize this run.""" t_stop: Time = Time.now() t_cost: float = (t_stop - self.t_start).value * 86400.0 self.logger.info(f"Total cost: {t_cost:.1f} sec") def clean_output(self): """Clean output directory.""" # self.clean_directory(self.dir_output) # clean in run command pass # file operations @staticmethod def mkdir(d): """Create a directory if it does not exist.""" if not os.path.exists(d): os.makedirs(d) @staticmethod def clean_directory(d): Loading @@ -158,19 +165,26 @@ class Pipeline: print("Failed to clean output directory!") print_directory_tree(d) def now(self): """Return ISOT format datetime using `astropy`.""" return Time.now().isot @staticmethod def filter_warnings(level: str = "ignore"): # Suppress all warnings warnings.filterwarnings(level) def move(file_src: str, file_dst: str) -> str: """Move file `file_src` to `file_dist`.""" return shutil.move(file_src, file_dst) @staticmethod def reset_warnings(self): """Reset warning filters.""" warnings.resetwarnings() def copy(file_src: str, file_dst: str) -> str: """Move file `file_src` to `file_dist`.""" return shutil.copy(file_src, file_dst) @property def summarize(self): """Summarize this run.""" t_stop: Time = Time.now() t_cost: float = (t_stop - self.t_start).value * 86400.0 self.logger.info(f"Total cost: {t_cost:.1f} sec") def clean_output(self): """Clean output directory.""" self.clean_directory(self.dir_output) def file(self, file_path): """Initialize File object.""" Loading @@ -180,25 +194,66 @@ class Pipeline: """Create new file in output directory.""" return os.path.join(self.dir_output, file_name) @staticmethod def move(file_src: str, file_dst: str): """Move file `file_src` to `file_dist`.""" shutil.move(file_src, file_dst) @staticmethod def copy(file_src: str, file_dst: str): """Move file `file_src` to `file_dist`.""" shutil.copy(file_src, file_dst) def download_oss_file(self, oss_file_path: str) -> str: def download_oss_file(self, oss_file_path: str, dir_dst: str = None) -> str: """Download an OSS file from OSS to output directory.""" local_file_path = os.path.join(self.dir_output, os.path.basename(oss_file_path)) if dir_dst is None: dir_dst = self.dir_output local_file_path = os.path.join(dir_dst, os.path.basename(oss_file_path)) csst_fs.s3_fs.get(oss_file_path, local_file_path, s3_options=s3_options) assert os.path.exists( local_file_path ), f"Failed to download {oss_file_path} to {local_file_path}" return local_file_path def abspath(self, file_path: str) -> str: """Return absolute path of `file_path`.""" if file_path.__contains__(":"): # it's an OSS file path assert self.use_oss, "USE_OSS must be True to use OSS file path!" # download OSS file to output directory local_file_path = self.download_oss_file(file_path) # return local file path return local_file_path else: # it's a NAS file path if file_path.startswith("CSST"): # DFS return os.path.join(self.dfs_root, file_path) else: # CCDS return os.path.join(self.ccds_root, file_path) def download_dfs_file(self, file_path: str, dir_dst: str = None): """Copy DFS file to output directory.""" if dir_dst is None: dir_dst = self.dir_output if self.use_oss: # download OSS file to dst directory return self.download_oss_file(file_path, dir_dst) else: # copy DFS file to dst directory local_file_path = os.path.join(dir_dst, os.path.basename(file_path)) self.copy(self.abspath(file_path), local_file_path) return local_file_path def copy_ccds_refs(self, refs: dict, dir_dst: str = None) -> dict: """Copy raw file from CCDS to output directory.""" if dir_dst is None: dir_dst = self.dir_output local_refs = {} for ref_name, ref_path in refs.items(): local_file_path = os.path.join(dir_dst, os.path.basename(ref_path)) local_refs[ref_name] = self.copy(self.abspath(ref_path), local_file_path) return local_refs # time operations def now(self) -> str: """Return ISOT format datetime using `astropy`.""" return Time.now().isot # call modules def call(self, func: Callable, *args: Any, **kwargs: Any): self.logger.info(f"=====================================================") t_start: Time = Time.now() Loading Loading @@ -249,6 +304,7 @@ class Pipeline: output=output, ) # retry operations @staticmethod def retry(*args, **kwargs): return retry(*args, **kwargs) Loading @@ -261,6 +317,7 @@ class Pipeline: else: raise ValueError("CCDS client not initialized!") # automatically log pipeline information def info(self): """Return pipeline information.""" self.logger.info(f"DOCKER_IMAGE={self.docker_image}") Loading