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.

This commit is contained in:
Kaivan Kamali
2020-08-25 16:42:33 -04:00
parent 08b9ea521c
commit 35d4e77015
+118 -19
View File
@@ -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=<copy|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 <copy|checksum>")
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)