From 35d4e770158022047a36b9c818299883f311bdaf Mon Sep 17 00:00:00 2001 From: Kaivan Kamali Date: Tue, 25 Aug 2020 16:42:33 -0400 Subject: [PATCH] Merged both scripts into one. The same script does the copy or checksum/delete depending on a new input parmater. Will delete the second script later. --- scripts/objectstore/copy_files_to_irods.py | 137 ++++++++++++++++++--- 1 file changed, 118 insertions(+), 19 deletions(-) diff --git a/scripts/objectstore/copy_files_to_irods.py b/scripts/objectstore/copy_files_to_irods.py index b70e42a8ba1..9ec69f29bf0 100644 --- a/scripts/objectstore/copy_files_to_irods.py +++ b/scripts/objectstore/copy_files_to_irods.py @@ -2,10 +2,14 @@ import getopt import json import os import uuid +import shutil import subprocess import sys import irods.keywords as kw +from irods.exception import CollectionDoesNotExist +from irods.exception import DataObjectDoesNotExist +from irods.exception import NetworkException from irods.session import iRODSSession from psycopg2 import connect @@ -13,7 +17,7 @@ from psycopg2 import connect from galaxy.util import directory_hash_id -def copy_files_to_irods(start_dataset_id, end_dataset_id, object_store_info_file, irods_info_file, db_connection_info_file): +def copy_files_to_irods(start_dataset_id, end_dataset_id, object_store_info_file, irods_info_file, db_connection_info_file, copy_or_checksum): conn = None session = None osi_keys = None @@ -83,6 +87,9 @@ def copy_files_to_irods(start_dataset_id, end_dataset_id, object_store_info_file AND id >= %s AND id <= %s AND object_store_id IN %s""" + update_sql_statement = """UPDATE dataset + SET object_store_id = %s + WHERE id = %s""" try: read_cursor = conn.cursor() args = ('ok', start_dataset_id, end_dataset_id, osi_keys) @@ -132,22 +139,96 @@ def copy_files_to_irods(start_dataset_id, end_dataset_id, object_store_info_file irods_folder_collection_path = os.path.join(irods_file_collection_path, "dataset_" + str(uuid_with_dash) + "_files") print("irods_folder_collection_path: ", irods_folder_collection_path) - # Create the collection - session.collections.create(irods_file_collection_path) - - # Add disk file to collection - options = {kw.DEST_RESC_NAME_KW: 'demoResc', kw.REG_CHKSUM_KW : ''} - session.data_objects.put(disk_file_path, irods_file_path, **options) - - if os.path.isdir(disk_folder_path): - disk_folder_path_all_files = disk_folder_path + "/*" - print("disk_folder_path_all_files: ", disk_folder_path_all_files) - + if copy_or_checksum == "copy": # Create the collection - session.collections.create(irods_folder_collection_path) + session.collections.create(irods_file_collection_path) - iput_command = "iput -rk " + disk_folder_path_all_files + " " + irods_folder_collection_path - subprocess.call(iput_command, shell=True) + # Add disk file to collection + options = {kw.DEST_RESC_NAME_KW: 'demoResc', kw.REG_CHKSUM_KW : ''} + session.data_objects.put(disk_file_path, irods_file_path, **options) + + if os.path.isdir(disk_folder_path): + disk_folder_path_all_files = disk_folder_path + "/*" + print("disk_folder_path_all_files: ", disk_folder_path_all_files) + + # Create the collection + session.collections.create(irods_folder_collection_path) + + iput_command = "iput -rk " + disk_folder_path_all_files + " " + irods_folder_collection_path + subprocess.call(iput_command, shell=True) + + if copy_or_checksum == "checksum": + # Calculate disk file checksum before uploading it to irods + # After disk file is uploaded to irods, we get the file checksum from irods and compare it with the calculated disk file checksum + # Note that disk file checksum is ASCII, whereas irods file checksum is Unicode + disk_file_checksum = get_file_checksum(disk_file_path) + print("disk_file_checksum: ", disk_file_checksum) + # Now get the file from irods + try: + obj = session.data_objects.get(irods_file_path) + print("obj.checksum: ", obj.checksum) + # obj.checksum is prepended with 'sha2:'. Remove that so we can compare it to disk file checksum + irods_file_checksum = obj.checksum[5:] + print("irods_file_checksum: ", irods_file_checksum) + if irods_file_checksum != disk_file_checksum: + print("irods file checksum {} does not match disk file checksum {}".format(irods_file_checksum, disk_file_checksum)) + continue + except (DataObjectDoesNotExist, CollectionDoesNotExist) as e: + print(e) + continue + except NetworkException as e: + print(e) + continue + + # Recursively verify that the checksum of all files in this folder matches that in irods + if os.path.isdir(disk_folder_path): + # Recursively traverse the files in this folder + for root, dirs, files in os.walk(disk_folder_path): + for file_name in files: + a_disk_file_path = os.path.join(root, file_name) + print(disk_file_path) + # Get checksum for disk file + a_disk_file_checksum = get_file_checksum(a_disk_file_path) + print("a_disk_file_checksum: ", a_disk_file_checksum) + + # Construct iords path for this disk file, so can get the file from irods, and compare its checksum with disk file checksum + # This is to extract the subfoler name for irods from the full disk path + sub_folder = root.replace(disk_folder_path + "/", "") + print("sub_folder: ", sub_folder) + an_irods_file_path = irods_folder_collection_path + "/" + sub_folder + "/" + file_name + print("an_irods_file_path: ", an_irods_file_path) + + # Now get the file from irods + try: + obj = session.data_objects.get(an_irods_file_path) + print("obj.checksum: ", obj.checksum) + # obj.checksum is prepended with 'sha2:'. Remove that so we can compare it to disk file checksum + an_irods_file_checksum = obj.checksum[5:] + print("an_irods_file_checksum: ", an_irods_file_checksum) + if an_irods_file_checksum != a_disk_file_checksum: + print("irods file checksum {} does not match disk file checksum {}".format(an_irods_file_checksum, a_disk_file_checksum)) + continue + except (DataObjectDoesNotExist, CollectionDoesNotExist) as e: + print(e) + continue + except NetworkException as e: + print(e) + continue + + # Delete file on disk + print("Removing directory " + disk_folder_path) + shutil.rmtree(disk_folder_path) + + # Update object store id + update_cursor = conn.cursor() + print("irods object store id: ", irods_info["object_store_id"]) + update_cursor.execute(update_sql_statement, (irods_info["object_store_id"], objectid)) + updated_rows = update_cursor.rowcount + print("updated_rows: ", updated_rows) + update_cursor.close() + + # Delete file on disk + os.remove(disk_file_path) except Exception as e: print(e) @@ -163,12 +244,22 @@ def copy_files_to_irods(start_dataset_id, end_dataset_id, object_store_info_file conn.close() +def get_file_checksum(disk_file_path): + checksum_cmd = "shasum -a 256 {} | xxd -r -p | base64".format(disk_file_path) + disk_file_checksum = subprocess.check_output(checksum_cmd, shell=True) + # remove '\n' from the end of disk_file_checksum + disk_file_checksum_len = len(disk_file_checksum) + disk_file_checksum_trimmed = disk_file_checksum[0:(disk_file_checksum_len - 1)] + # Return Unicode string + return disk_file_checksum_trimmed.decode("utf-8") + + def print_help_msg(): print("\nLong form input parameter specification:") - print("copy_files_to_irods --start_dataset_id=2 --end_dataset_id=3 --object_store_info_file=object_store_info.json --irods_info_file=irods_info_file.json --db_connection_info_file=db_connection_info_file.json") + print("copy_files_to_irods --start_dataset_id=2 --end_dataset_id=3 --object_store_info_file=object_store_info.json --irods_info_file=irods_info_file.json --db_connection_info_file=db_connection_info_file.json --cop_or_checksum=") print("\nOR") print("\nShort form input parameter specification:") - print("copy_files_to_irods -s 2 -e 3 -o object_store_info.json -i irods_info_file.json -d db_connection_info_file.json") + print("copy_files_to_irods -s 2 -e 3 -o object_store_info.json -i irods_info_file.json -d db_connection_info_file.json -c ") print("\n") @@ -178,9 +269,10 @@ if __name__ == '__main__': object_store_info_file = None irods_info_file = None db_connection_info_file = None + copy_or_checksum = None try: - opts, args = getopt.getopt(sys.argv[1:], "hs:e:o:i:d:", ["start_dataset_id=", "end_dataset_id=", "object_store_info_file=", "irods_info_file=", "db_connection_info_file="]) + opts, args = getopt.getopt(sys.argv[1:], "hs:e:o:i:d:c:", ["start_dataset_id=", "end_dataset_id=", "object_store_info_file=", "irods_info_file=", "db_connection_info_file=", "copy_or_checksum="]) except getopt.GetoptError: print_help_msg() sys.exit(2) @@ -199,10 +291,17 @@ if __name__ == '__main__': irods_info_file = arg elif opt in ("-d", "--db_connection-info-file"): db_connection_info_file = arg + elif opt in ("-c", "--copy_or_checksum"): + copy_or_checksum = arg if start_dataset_id is None or end_dataset_id is None or object_store_info_file is None or irods_info_file is None or db_connection_info_file is None: print("Did not specify one of the required input parameters!") print_help_msg() sys.exit(2) - copy_files_to_irods(start_dataset_id=start_dataset_id, end_dataset_id=end_dataset_id, object_store_info_file=object_store_info_file, irods_info_file=irods_info_file, db_connection_info_file=db_connection_info_file) + if copy_or_checksum != "copy" and copy_or_checksum != "checksum": + print("Did not specify correct value for copy_or_checksum input parameter!") + print_help_msg() + sys.exit(2) + + copy_files_to_irods(start_dataset_id=start_dataset_id, end_dataset_id=end_dataset_id, object_store_info_file=object_store_info_file, irods_info_file=irods_info_file, db_connection_info_file=db_connection_info_file, copy_or_checksum=copy_or_checksum)