From 1e9764815e45fd333b20e09a4431a2fd04883863 Mon Sep 17 00:00:00 2001 From: dert1129 Date: Fri, 17 Jul 2026 12:06:30 -0400 Subject: [PATCH 1/2] refactor copy files method. if there is a top level directory, remove it but keep the structure inside --- data_management/services/dlu_filesystem.py | 344 +++++++++++++++++---- 1 file changed, 288 insertions(+), 56 deletions(-) diff --git a/data_management/services/dlu_filesystem.py b/data_management/services/dlu_filesystem.py index f1b9d44..98ad2f1 100644 --- a/data_management/services/dlu_filesystem.py +++ b/data_management/services/dlu_filesystem.py @@ -135,68 +135,300 @@ def rename_and_move_files(self, file_list: list[DLUFile], slide_name_map, packag checksum=calculate_checksum(dest_file), size=os.path.getsize(dest_file)) dluFiles.append(file) return dluFiles + + def copy_directory_contents(src_dir: str, dst_dir: str) -> int: + if not os.path.isdir(src_dir): + raise FileNotFoundError(f"Source directory does not exist: {src_dir}") + + logger.info("Copying contents of %s into %s", src_dir, dst_dir) + + os.makedirs(dst_dir, exist_ok=True) - def copy_files(self, package_id: str, file_list: list[DLUFile], preserve_path: bool = False, no_src_package: bool = False): files_copied = 0 - source_wd = os.getcwd() - dest_package_directory = os.path.join(self.dlu_data_directory, self.dlu_package_dir_prefix + package_id) - if os.path.exists(dest_package_directory): - shutil.rmtree(dest_package_directory) + + for root, dirnames, filenames in os.walk(src_dir): + relative_root = os.path.relpath(root, src_dir) + + if relative_root == ".": + target_root = dst_dir + else: + target_root = os.path.join(dst_dir, relative_root) + + os.makedirs(target_root, exist_ok=True) + + # Create directories even if they are empty. + for dirname in dirnames: + target_dir = os.path.join(target_root, dirname) + os.makedirs(target_dir, exist_ok=True) + + for filename in filenames: + src_file = os.path.join(root, filename) + dst_file = os.path.join(target_root, filename) + + if os.path.exists(dst_file): + logger.warning("%s already exists. Skipping.", dst_file) + continue + + logger.info("Copying file %s to %s", src_file, dst_file) + shutil.copy2(src_file, dst_file) + files_copied += 1 + + return files_copied + + def copy_files(self,package_id: str, file_list: list[DLUFile], preserve_path: bool = False, no_src_package: bool = False): + files_copied = 0 + + base_dest_package_directory = os.path.join( + self.dlu_data_directory, + self.dlu_package_dir_prefix + package_id, + ) + + if os.path.exists(base_dest_package_directory): + logger.info( + "Removing existing destination directory %s", + base_dest_package_directory, + ) + shutil.rmtree(base_dest_package_directory) + + # Used to prevent copying the same wrapper directory repeatedly if file_list has multiple files from the same package. + copied_wrapper_directories = set() + + def count_files(path: str) -> int: + """ + Count regular files under path. + If path is a file, returns 1. + If path is a directory, recursively counts files. + """ + if os.path.isfile(path): + return 1 + + total = 0 + + for _, _, filenames in os.walk(path): + total += len(filenames) + + return total + + def copy_path(src_path: str, dst_path: str) -> int: + """ + Copy a file or directory from src_path to dst_path. + + Returns the number of files copied. + """ + + logger.info("Preparing to copy from %s to %s", src_path, dst_path) + + if not os.path.exists(src_path): + raise FileNotFoundError(f"Source path does not exist: {src_path}") + + dst_parent = os.path.dirname(dst_path) + + if dst_parent: + os.makedirs(dst_parent, exist_ok=True) + + if os.path.exists(dst_path): + logger.warning("%s already exists. Skipping.", dst_path) + return 0 + + try: + if os.path.isdir(src_path): + logger.info("Copying directory %s to %s", src_path, dst_path) + shutil.copytree(src_path, dst_path) + return count_files(src_path) + + if os.path.isfile(src_path): + logger.info("Copying file %s to %s", src_path, dst_path) + shutil.copy2(src_path, dst_path) + return 1 + + raise FileNotFoundError( + f"Source path exists but is neither a regular file nor directory: " + f"{src_path}" + ) + + except FileNotFoundError: + logger.exception( + "Source disappeared or destination parent was missing during copy." + ) + raise + + def copy_directory_contents(src_dir: str, dst_dir: str) -> int: + + if not os.path.isdir(src_dir): + raise FileNotFoundError(f"Source directory does not exist: {src_dir}") + + logger.info("Copying contents of %s into %s", src_dir, dst_dir) + + os.makedirs(dst_dir, exist_ok=True) + + copied_count = 0 + + for root, dirnames, filenames in os.walk(src_dir): + relative_root = os.path.relpath(root, src_dir) + + if relative_root == ".": + target_root = dst_dir + else: + target_root = os.path.join(dst_dir, relative_root) + + os.makedirs(target_root, exist_ok=True) + + # Create directories even if they are empty. + for dirname in dirnames: + target_dir = os.path.join(target_root, dirname) + os.makedirs(target_dir, exist_ok=True) + + for filename in filenames: + src_file = os.path.join(root, filename) + dst_file = os.path.join(target_root, filename) + + if os.path.exists(dst_file): + logger.warning("%s already exists. Skipping.", dst_file) + continue + + logger.info("Copying file %s to %s", src_file, dst_file) + shutil.copy2(src_file, dst_file) + copied_count += 1 + + return copied_count + for file in file_list: + source_package_directory = os.path.join( + self.globus_data_directory, + self.globus_dir_prefix + ("" if no_src_package else package_id), + ) + + dest_package_directory = base_dest_package_directory - source_package_directory = self.globus_data_directory + '/' + self.globus_dir_prefix - # I.e. isn't a bulk upload that doesn't already have a package ID. - logger.info(source_package_directory) - if not no_src_package: - source_package_directory = source_package_directory + package_id - if file.path and os.path.isdir(file.path): - source_package_directory = os.path.join(source_package_directory, file.path) if preserve_path: - dest_package_directory = os.path.join(dest_package_directory, - file.get_short_path()) - - subdirs = [os.path.join(source_package_directory, o) - for o in os.listdir(source_package_directory) - if os.path.isdir(os.path.join(source_package_directory, o))] - dir = "".join(subdirs) - if len(os.listdir(source_package_directory)) == 1 and os.path.isdir(source_package_directory) and os.path.isdir(dir): - os.chdir(dir) - allfiles = os.listdir(dir) - for f in allfiles: - src_path = os.path.join(dir, f) - dst_path = os.path.join(dest_package_directory, f) - if not os.path.isdir(dest_package_directory): - os.mkdir(dest_package_directory) - if os.path.isfile(f): - logger.info("Copying file " + f + " to " + dst_path) - shutil.copy(src_path, dst_path) - files_copied += 1 + short_path = file.get_short_path() + + if short_path: + dest_package_directory = os.path.join( + dest_package_directory, + short_path, + ) + + logger.info("Base source package directory: %s", source_package_directory) + logger.info("Destination package directory: %s", dest_package_directory) + logger.info("File path: %s", getattr(file, "path", None)) + logger.info("Short filename: %s", file.get_short_filename()) + + if not os.path.isdir(source_package_directory): + raise FileNotFoundError( + f"Source package directory does not exist or is not a directory: " + f"{source_package_directory}" + ) + + os.makedirs(dest_package_directory, exist_ok=True) + # Special top-level wrapper directory behavior. + # Copy the contents of top_level_dir, but do not copy top_level_dir itself. + + try: + top_level_items = os.listdir(source_package_directory) + except FileNotFoundError: + raise FileNotFoundError( + f"Cannot list source package directory because it does not exist: " + f"{source_package_directory}" + ) + + if len(top_level_items) == 1: + only_item_name = top_level_items[0] + only_item_path = os.path.join(source_package_directory, only_item_name) + + if os.path.isdir(only_item_path): + wrapper_key = ( + os.path.abspath(only_item_path), + os.path.abspath(dest_package_directory), + ) + + if wrapper_key not in copied_wrapper_directories: + logger.info( + "Source package directory contains a single wrapper directory. " + "Copying contents of %s into %s", + only_item_path, + dest_package_directory, + ) + + files_copied += copy_directory_contents( + src_dir=only_item_path, + dst_dir=dest_package_directory, + ) + + copied_wrapper_directories.add(wrapper_key) + else: - logger.info("Copying directory " + src_path) - files_copied += 1 - shutil.copytree(src_path, dst_path) - os.chdir(source_wd) - - if not os.path.exists(dest_package_directory): - logger.info("Creating directory " + dest_package_directory) - os.makedirs(dest_package_directory, exist_ok=True) - source_file = os.path.join(source_package_directory, file.get_short_filename()) - dest_file = os.path.join(dest_package_directory, file.get_short_filename()) - - if not os.path.exists(dest_file): - if os.path.isdir(source_file): - logger.info("Copying directory to " + dest_file) - shutil.copytree(source_file, dest_file) - elif os.path.isfile(source_file): - logger.info("Copying file to " + dest_file) - shutil.copy(source_file, dest_file) - else: - source_file = os.path.join(source_package_directory, file.path) - logger.info("Copying file to " + dest_file) - shutil.copy(source_file, dest_file) - files_copied = files_copied + 1 - else: - logger.warning(dest_file + " already exists. Skipping.") + logger.info( + "Wrapper directory %s has already been copied to %s. Skipping.", + only_item_path, + dest_package_directory, + ) + + # Since the entire wrapper directory contents were copied, + # do not continue into individual file copy logic for this file. + continue + + source_directory_for_file = source_package_directory + + # If file.path points to a directory under the source package directory, + # use that as the source directory for this file. + if file.path: + candidate_source_directory = os.path.join( + source_package_directory, + file.path, + ) + + if os.path.isdir(candidate_source_directory): + source_directory_for_file = candidate_source_directory + + short_filename = file.get_short_filename() + + source_file = os.path.join( + source_directory_for_file, + short_filename, + ) + + dest_file = os.path.join( + dest_package_directory, + short_filename, + ) + + if os.path.exists(source_file): + files_copied += copy_path(source_file, dest_file) + continue + + # Fallback: if source_file was not found, try file.path directly. + fallback_source_file = None + + if file.path: + fallback_source_file = os.path.join( + source_package_directory, + file.path, + ) + + if os.path.exists(fallback_source_file): + fallback_dest_file = os.path.join( + dest_package_directory, + os.path.basename(file.path.rstrip(os.sep)), + ) + + files_copied += copy_path( + fallback_source_file, + fallback_dest_file, + ) + + continue + + raise FileNotFoundError( + "Could not find source file or directory. Tried:\n" + f" source_file={source_file}\n" + f" fallback_source_file={fallback_source_file}\n" + f" source_package_directory={source_package_directory}\n" + f" source_directory_for_file={source_directory_for_file}\n" + f" file.path={getattr(file, 'path', None)}\n" + f" short_filename={short_filename}" + ) + return files_copied def validate_package_directories(self, package_id: str): From 752066121b3dd57f31e78a62f24f705a368b97f6 Mon Sep 17 00:00:00 2001 From: dert1129 Date: Fri, 17 Jul 2026 15:13:31 -0400 Subject: [PATCH 2/2] try and make directory removal a bit safer --- data_management/services/dlu_filesystem.py | 379 ++++++++++++++------- 1 file changed, 260 insertions(+), 119 deletions(-) diff --git a/data_management/services/dlu_filesystem.py b/data_management/services/dlu_filesystem.py index 98ad2f1..63702cc 100644 --- a/data_management/services/dlu_filesystem.py +++ b/data_management/services/dlu_filesystem.py @@ -9,6 +9,7 @@ from zarr_checksum.generators import yield_files_local from mmap import mmap, ACCESS_READ import subprocess +import tempfile logger = logging.getLogger("DLUFilesystem") logger.setLevel(logging.INFO) @@ -175,30 +176,60 @@ def copy_directory_contents(src_dir: str, dst_dir: str) -> int: return files_copied - def copy_files(self,package_id: str, file_list: list[DLUFile], preserve_path: bool = False, no_src_package: bool = False): + def copy_files( self, package_id: str, file_list: list[DLUFile], preserve_path: bool = False, no_src_package: bool = False): files_copied = 0 - base_dest_package_directory = os.path.join( + final_dest_package_directory = os.path.join( self.dlu_data_directory, self.dlu_package_dir_prefix + package_id, ) - if os.path.exists(base_dest_package_directory): - logger.info( - "Removing existing destination directory %s", - base_dest_package_directory, + dest_parent_directory = os.path.dirname(final_dest_package_directory) + dest_package_basename = os.path.basename(final_dest_package_directory) + + if not file_list: + logger.warning( + "No files provided for package %s. Existing package will not be modified.", + package_id, + ) + return 0 + + if not os.path.isdir(dest_parent_directory): + raise FileNotFoundError( + f"Destination parent directory does not exist: {dest_parent_directory}" ) - shutil.rmtree(base_dest_package_directory) - # Used to prevent copying the same wrapper directory repeatedly if file_list has multiple files from the same package. + if not os.access(dest_parent_directory, os.W_OK | os.X_OK): + raise PermissionError( + f"Process does not have permission to write to destination parent directory: " + f"{dest_parent_directory}" + ) + + if os.path.exists(final_dest_package_directory) and not os.path.isdir( + final_dest_package_directory + ): + raise NotADirectoryError( + f"Destination package path exists but is not a directory: " + f"{final_dest_package_directory}" + ) + + temp_dest_package_directory = tempfile.mkdtemp( + prefix=f".{dest_package_basename}.tmp-", + dir=dest_parent_directory, + ) + + backup_dest_package_directory = None + + logger.info( + "Building replacement package %s in temporary directory %s", + package_id, + temp_dest_package_directory, + ) + copied_wrapper_directories = set() def count_files(path: str) -> int: - """ - Count regular files under path. - If path is a file, returns 1. - If path is a directory, recursively counts files. - """ + if os.path.isfile(path): return 1 @@ -210,11 +241,6 @@ def count_files(path: str) -> int: return total def copy_path(src_path: str, dst_path: str) -> int: - """ - Copy a file or directory from src_path to dst_path. - - Returns the number of files copied. - """ logger.info("Preparing to copy from %s to %s", src_path, dst_path) @@ -253,7 +279,6 @@ def copy_path(src_path: str, dst_path: str) -> int: raise def copy_directory_contents(src_dir: str, dst_dir: str) -> int: - if not os.path.isdir(src_dir): raise FileNotFoundError(f"Source directory does not exist: {src_dir}") @@ -292,145 +317,261 @@ def copy_directory_contents(src_dir: str, dst_dir: str) -> int: return copied_count - for file in file_list: - source_package_directory = os.path.join( - self.globus_data_directory, - self.globus_dir_prefix + ("" if no_src_package else package_id), - ) + def replace_destination_with_temp(): - dest_package_directory = base_dest_package_directory - - if preserve_path: - short_path = file.get_short_path() + nonlocal backup_dest_package_directory - if short_path: - dest_package_directory = os.path.join( - dest_package_directory, - short_path, - ) + if os.path.exists(final_dest_package_directory): + backup_dest_package_directory = tempfile.mkdtemp( + prefix=f".{dest_package_basename}.backup-", + dir=dest_parent_directory, + ) - logger.info("Base source package directory: %s", source_package_directory) - logger.info("Destination package directory: %s", dest_package_directory) - logger.info("File path: %s", getattr(file, "path", None)) - logger.info("Short filename: %s", file.get_short_filename()) + # tempfile.mkdtemp creates the backup directory. Remove it so + # os.rename can move the existing package to that path. + os.rmdir(backup_dest_package_directory) - if not os.path.isdir(source_package_directory): - raise FileNotFoundError( - f"Source package directory does not exist or is not a directory: " - f"{source_package_directory}" + logger.info( + "Moving existing destination %s to backup %s", + final_dest_package_directory, + backup_dest_package_directory, ) - os.makedirs(dest_package_directory, exist_ok=True) - # Special top-level wrapper directory behavior. - # Copy the contents of top_level_dir, but do not copy top_level_dir itself. + os.rename( + final_dest_package_directory, + backup_dest_package_directory, + ) try: - top_level_items = os.listdir(source_package_directory) - except FileNotFoundError: - raise FileNotFoundError( - f"Cannot list source package directory because it does not exist: " - f"{source_package_directory}" + logger.info( + "Moving temporary package %s into final destination %s", + temp_dest_package_directory, + final_dest_package_directory, + ) + + os.rename( + temp_dest_package_directory, + final_dest_package_directory, + ) + + except Exception: + logger.exception( + "Failed to move temporary package into final destination." ) - if len(top_level_items) == 1: - only_item_name = top_level_items[0] - only_item_path = os.path.join(source_package_directory, only_item_name) + if ( + backup_dest_package_directory + and os.path.exists(backup_dest_package_directory) + and not os.path.exists(final_dest_package_directory) + ): + logger.warning( + "Attempting to restore backup package from %s to %s", + backup_dest_package_directory, + final_dest_package_directory, + ) + + os.rename( + backup_dest_package_directory, + final_dest_package_directory, + ) + + raise - if os.path.isdir(only_item_path): - wrapper_key = ( - os.path.abspath(only_item_path), - os.path.abspath(dest_package_directory), + if ( + backup_dest_package_directory + and os.path.exists(backup_dest_package_directory) + ): + try: + logger.info( + "Removing old backup package directory %s", + backup_dest_package_directory, + ) + shutil.rmtree(backup_dest_package_directory) + except Exception: + logger.exception( + "Package replacement succeeded, but backup cleanup failed: %s", + backup_dest_package_directory, ) - if wrapper_key not in copied_wrapper_directories: - logger.info( - "Source package directory contains a single wrapper directory. " - "Copying contents of %s into %s", - only_item_path, + try: + for file in file_list: + source_package_directory = os.path.join( + self.globus_data_directory, + self.globus_dir_prefix + ("" if no_src_package else package_id), + ) + + # This is the base temporary package directory. + # All copying happens here first, not directly into the final destination. + dest_package_directory = temp_dest_package_directory + + if preserve_path: + short_path = file.get_short_path() + + if short_path: + dest_package_directory = os.path.join( dest_package_directory, + short_path, ) - files_copied += copy_directory_contents( - src_dir=only_item_path, - dst_dir=dest_package_directory, - ) + logger.info("Base source package directory: %s", source_package_directory) + logger.info("Temporary package base directory: %s", temp_dest_package_directory) + logger.info("Per-file destination directory: %s", dest_package_directory) + logger.info("File path: %s", getattr(file, "path", None)) + logger.info("Short filename: %s", file.get_short_filename()) - copied_wrapper_directories.add(wrapper_key) + if not os.path.isdir(source_package_directory): + raise FileNotFoundError( + f"Source package directory does not exist or is not a directory: " + f"{source_package_directory}" + ) - else: - logger.info( - "Wrapper directory %s has already been copied to %s. Skipping.", - only_item_path, - dest_package_directory, + os.makedirs(dest_package_directory, exist_ok=True) + + # If source_package_directory contains exactly one item and that + # item is a directory, treat it as a top-level wrapper directory. + # Important: + # Even with preserve_path=True, wrapper contents are copied into + # temp_dest_package_directory, not dest_package_directory. + + try: + top_level_items = os.listdir(source_package_directory) + except FileNotFoundError: + raise FileNotFoundError( + f"Cannot list source package directory because it does not exist: " + f"{source_package_directory}" + ) + + if len(top_level_items) == 1: + only_item_name = top_level_items[0] + only_item_path = os.path.join( + source_package_directory, + only_item_name, + ) + + if os.path.isdir(only_item_path): + wrapper_destination = temp_dest_package_directory + + wrapper_key = ( + os.path.abspath(only_item_path), + os.path.abspath(wrapper_destination), ) - # Since the entire wrapper directory contents were copied, - # do not continue into individual file copy logic for this file. - continue + if wrapper_key not in copied_wrapper_directories: + logger.info( + "Source package directory contains a single wrapper directory. " + "Copying contents of %s into base package destination %s", + only_item_path, + wrapper_destination, + ) - source_directory_for_file = source_package_directory + files_copied += copy_directory_contents( + src_dir=only_item_path, + dst_dir=wrapper_destination, + ) - # If file.path points to a directory under the source package directory, - # use that as the source directory for this file. - if file.path: - candidate_source_directory = os.path.join( - source_package_directory, - file.path, - ) + copied_wrapper_directories.add(wrapper_key) - if os.path.isdir(candidate_source_directory): - source_directory_for_file = candidate_source_directory + else: + logger.info( + "Wrapper directory %s has already been copied to %s. Skipping.", + only_item_path, + wrapper_destination, + ) - short_filename = file.get_short_filename() + continue - source_file = os.path.join( - source_directory_for_file, - short_filename, - ) + source_directory_for_file = source_package_directory - dest_file = os.path.join( - dest_package_directory, - short_filename, - ) + if file.path: + candidate_source_directory = os.path.join( + source_package_directory, + file.path, + ) - if os.path.exists(source_file): - files_copied += copy_path(source_file, dest_file) - continue + if os.path.isdir(candidate_source_directory): + source_directory_for_file = candidate_source_directory - # Fallback: if source_file was not found, try file.path directly. - fallback_source_file = None + short_filename = file.get_short_filename() - if file.path: - fallback_source_file = os.path.join( - source_package_directory, - file.path, + source_file = os.path.join( + source_directory_for_file, + short_filename, ) - if os.path.exists(fallback_source_file): - fallback_dest_file = os.path.join( - dest_package_directory, - os.path.basename(file.path.rstrip(os.sep)), - ) + dest_file = os.path.join( + dest_package_directory, + short_filename, + ) + if os.path.exists(source_file): files_copied += copy_path( - fallback_source_file, - fallback_dest_file, + src_path=source_file, + dst_path=dest_file, ) - continue - raise FileNotFoundError( - "Could not find source file or directory. Tried:\n" - f" source_file={source_file}\n" - f" fallback_source_file={fallback_source_file}\n" - f" source_package_directory={source_package_directory}\n" - f" source_directory_for_file={source_directory_for_file}\n" - f" file.path={getattr(file, 'path', None)}\n" - f" short_filename={short_filename}" + fallback_source_file = None + + if file.path: + fallback_source_file = os.path.join( + source_package_directory, + file.path, + ) + + if os.path.exists(fallback_source_file): + fallback_dest_file = os.path.join( + dest_package_directory, + os.path.basename(file.path.rstrip(os.sep)), + ) + + files_copied += copy_path( + src_path=fallback_source_file, + dst_path=fallback_dest_file, + ) + + continue + + raise FileNotFoundError( + "Could not find source file or directory. Tried:\n" + f" source_file={source_file}\n" + f" fallback_source_file={fallback_source_file}\n" + f" source_package_directory={source_package_directory}\n" + f" source_directory_for_file={source_directory_for_file}\n" + f" file.path={getattr(file, 'path', None)}\n" + f" short_filename={short_filename}" + ) + + # The entire temporary package was built successfully. + # Only now replace the existing package. + replace_destination_with_temp() + + logger.info( + "Successfully copied %s files for package %s into %s", + files_copied, + package_id, + final_dest_package_directory, ) - return files_copied + return files_copied + except Exception: + logger.exception( + "Failed to build replacement package for %s. Existing package was not replaced.", + package_id, + ) + + # If failure happens before final replacement, remove temp dir. + if os.path.exists(temp_dest_package_directory): + logger.info( + "Removing failed temporary package directory %s", + temp_dest_package_directory, + ) + + shutil.rmtree( + temp_dest_package_directory, + ignore_errors=True, + ) + raise def validate_package_directories(self, package_id: str): source_package_directory = self.globus_data_directory + '/' + self.globus_dir_prefix + package_id source_directory_info = DirectoryInfo(source_package_directory, False)