diff --git a/pyproject.toml b/pyproject.toml index dd0fc4e..69dfa2b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "rda_python_common" -version = "2.0.16" +version = "2.0.17" authors = [ { name="Zaihua Ji", email="zji@ucar.edu" }, ] diff --git a/src/rda_python_common/pg_file.py b/src/rda_python_common/pg_file.py index eac0111..88567ad 100644 --- a/src/rda_python_common/pg_file.py +++ b/src/rda_python_common/pg_file.py @@ -73,7 +73,8 @@ def __init__(self): 'G': self.PGLOG['GPFSNAME'], 'O': self.OHOST, 'B': self.BHOST, - 'D': self.DHOST + 'D': self.DHOST, + } self.DPATHS = { 'G': self.PGLOG['DSSDATA'], @@ -88,7 +89,7 @@ def __init__(self): 'F': 'FAILED', } self.QPOINTS = { - 'L': 'gdex-glade', + 'L': 'gdex-glade', # or gdex-lustre 'B': 'gdex-quasar', 'D': 'gdex-quasar-drdata' } @@ -252,7 +253,7 @@ def local_copy_object(self, tofile, fromfile, bucket = None, meta = None, logact if 'user' not in meta: meta['user'] = self.PGLOG['CURUID'] if 'group' not in meta: meta['group'] = self.PGLOG['GDEXGRP'] uinfo = json.dumps(meta) - finfo = self.check_local_file(fromfile, 0, logact) + finfo = self.check_local_file(fromfile, 0, logact|self.PFSIZE) if not finfo: if finfo != None: return self.FAILURE return self.lmsg(fromfile, "{} to copy to {}-{}".format(self.PGLOG['MISSFILE'], self.OHOST, tofile), logact) @@ -260,7 +261,8 @@ def local_copy_object(self, tofile, fromfile, bucket = None, meta = None, logact tinfo = self.check_object_file(tofile, bucket, 0, logact) if tinfo and tinfo['data_size'] > 0: return self.pglog("{}-{}-{}: file exists already".format(self.OHOST, bucket, tofile), logact) - cmd = "{} ul -lf {} -b {} -k {} -md '{}'".format(self.OBJCTCMD, fromfile, bucket, tofile, uinfo) + ocmd = self.valid_command(self.OBJCTCMD, logact) + cmd = "{} ul -lf {} -b {} -k {} -md '{}'".format(ocmd, fromfile, bucket, tofile, uinfo) for loop in range(2): buf = self.pgsystem(cmd, logact, self.CMDBTH) tinfo = self.check_object_file(tofile, bucket, 0, logact) @@ -293,7 +295,8 @@ def quasar_multiple_trasnfer(self, tofiles, fromfiles, topoint, frompoint, logac destination_endpoint = topoint label = f"{self.ENDPOINTS[frompoint]} to {self.ENDPOINTS[topoint]} {action}" verify_checksum = True - cmd = f'{self.BACKCMD} {action} -se {source_endpoint} -de {destination_endpoint} --label "{label}"' + bcmd = self.valid_command(self.BACKCMD, logact) + cmd = f'{bcmd} {action} -se {source_endpoint} -de {destination_endpoint} --label "{label}"' if verify_checksum: cmd += ' -vc' cmd += ' --batch -' @@ -322,14 +325,14 @@ def endpoint_copy_endpoint(self, tofile, fromfile, topoint, frompoint, logact = if tinfo and tinfo['data_size'] > 0: return self.pglog("{}-{}: file exists already".format(topoint, tofile), logact) action = 'transfer' - cmd = f'{self.BACKCMD} {action} -se {frompoint} -de {topoint} -sf {fromfile} -df {tofile} -vc' + bcmd = self.valid_command(self.BACKCMD, logact) + cmd = f'{bcmd} {action} -se {frompoint} -de {topoint} -sf {fromfile} -df {tofile} -vc' task = self.submit_globus_task(cmd, topoint, logact) if task['stat'] == 'S': ret = self.SUCCESS elif task['stat'] == 'A': self.TASKIDS["{}-{}".format(topoint, tofile)] = task['id'] ret = self.FINISH - return ret # submit a globus task and return a task id @@ -370,7 +373,8 @@ def check_globus_status(self, taskid, endpoint = None, logact = 0): if not taskid: return ret if not endpoint: endpoint = self.PGLOG['BACKUPEP'] mp = r'Status:\s+({})'.format('|'.join(self.QSTATS.values())) - cmd = f"{self.BACKCMD} get-task {taskid}" + bcmd = self.valid_command(self.BACKCMD, logact) + cmd = f"{bcmd} get-task {taskid}" astats = ['OK', 'Queued'] for loop in range(2): buf = self.pgsystem(cmd, logact, self.CMDRET) @@ -386,7 +390,7 @@ def check_globus_status(self, taskid, endpoint = None, logact = 0): if logact&self.NOWAIT: errmsg = "{}: Cancel Task due to {}:\n{}".format(taskid, detail, buf) self.errlog(errmsg, 'B', 1, logact) - ccmd = f"{self.BACKCMD} cancel-task {taskid}" + ccmd = f"{bcmd} cancel-task {taskid}" self.pgsystem(ccmd, logact, 7) else: time.sleep(self.PGSIG['ETIME']) @@ -498,7 +502,8 @@ def object_copy_local(self, tofile, fromfile, bucket = None, logact = 0): if not finfo: if finfo != None: return ret return self.lmsg(fromfile, "{}-{} to copy to {}".format(self.OHOST, self.PGLOG['MISSFILE'], tofile), logact) - cmd = "{} go -k {} -b {}".format(self.OBJCTCMD, fromfile, bucket) + ocmd = self.valid_command(self.OBJCTCMD, logact) + cmd = "{} go -k {} -b {}".format(ocmd, fromfile, bucket) fromname = op.basename(fromfile) toname = op.basename(tofile) if toname == tofile: @@ -509,7 +514,7 @@ def object_copy_local(self, tofile, fromfile, bucket = None, logact = 0): loop = reset = 0 while (loop-reset) < 2: buf = self.pgsystem(cmd, logact, self.CMDBTH) - info = self.check_local_file(fromname, 143, logact) # 1+2+4+8+128 + info = self.check_local_file(fromname, 143, logact|self.PFSIZE) # 1+2+4+8+128 if info: if info['data_size'] == finfo['data_size']: self.set_local_mode(fromfile, info['isfile'], 0, info['mode'], info['logname'], logact) @@ -603,12 +608,13 @@ def delete_remote_file(self, file, host, logact = 0): # Delete a file on object store def delete_object_file(self, file, bucket = None, logact = 0): if not bucket: bucket = self.PGLOG['OBJCTBKT'] + ocmd = self.valid_command(self.OBJCTCMD, logact) for loop in range(2): list = self.object_glob(file, bucket, 0, logact) if not list: return self.FAILURE errmsg = None for key in list: - cmd = "{} dl {} -b {}".format(self.OBJCTCMD, key, bucket) + cmd = "{} dl {} -b {}".format(ocmd, key, bucket) if not self.pgsystem(cmd, logact, self.CMDERR): errmsg = self.PGLOG['SYSERR'] break @@ -622,7 +628,8 @@ def delete_backup_file(self, file, endpoint = None, logact = 0): if not endpoint: endpoint = self.PGLOG['BACKUPEP'] info = self.check_backup_file(file, endpoint, 0, logact) if not info: return self.FAILURE - cmd = f"{self.BACKCMD} delete -ep {endpoint} -tf {file}" + bcmd = self.valid_command(self.BACKCMD, logact) + cmd = f"{bcmd} delete -ep {endpoint} -tf {file}" task = self.submit_globus_task(cmd, endpoint, logact) if task['stat'] == 'S': return self.SUCCESS @@ -775,8 +782,9 @@ def move_object_file(self, tofile, fromfile, tobucket, frombucket, logact = 0): return self.errlog("{}-{}: Object File exists, cannot move {}-{} to it".format(tobucket, tofile, frombucket, fromfile), 'R', 1, logact) elif tinfo != None: return self.FAILURE - cmd = "{} mv -b {} -db {} -k {} -dk {}".format(self.OBJCTCMD, frombucket, tobucket, fromfile, tofile) - ucmd = "{} gm -k {} -b {}".format(self.OBJCTCMD, fromfile, frombucket) + ocmd = self.valid_command(self.OBJCTCMD, logact) + cmd = "{} mv -b {} -db {} -k {} -dk {}".format(ocmd, frombucket, tobucket, fromfile, tofile) + ucmd = "{} gm -k {} -b {}".format(ocmd, fromfile, frombucket) ubuf = self.pgsystem(ucmd, self.LOGWRN, self.CMDRET) if ubuf and re.match(r'^\{', ubuf): cmd += " -md '{}'".format(ubuf) for loop in range(2): @@ -808,7 +816,8 @@ def move_object_path(self, topath, frompath, tobucket, frombucket, logact = 0): return self.SUCCESS else: return self.errlog("{}-{}: {} to move".format(frombucket, frompath, self.PGLOG['MISSFILE']), 'R', 1, logact) - cmd = "{} mv -b {} -db {} -k {} -dk {}".format(self.OBJCTCMD, frombucket, tobucket, frompath, topath) + ocmd = self.valid_command(self.OBJCTCMD, logact) + cmd = "{} mv -b {} -db {} -k {} -dk {}".format(ocmd, frombucket, tobucket, frompath, topath) for loop in range(2): buf = self.pgsystem(cmd, logact, self.CMDBTH) fcnt = self.check_object_path(frompath, frombucket, logact) @@ -837,7 +846,8 @@ def move_backup_file(self, tofile, fromfile, endpoint = None, logact = 0): return self.errlog("{}: File exists, cannot move {} to it".format(tofile, fromfile), 'B', 1, logact) elif tinfo != None: return ret - cmd = f"{self.BACKCMD} rename -ep {endpoint} --old-path {fromfile} --new-path {tofile}" + bcmd = self.valid_command(self.BACKCMD, logact) + cmd = f"{bcmd} rename -ep {endpoint} --old-path {fromfile} --new-path {tofile}" loop = 0 while loop < 2: buf = self.pgsystem(cmd, logact, self.CMDRET) @@ -935,7 +945,8 @@ def make_one_backup_directory(self, dir, odir, endpoint = None, logact = 0): return self.FAILURE if not odir: odir = dir if not self.make_one_backup_directory(op.dirname(dir), odir, endpoint, logact): return self.FAILURE - cmd = f"{self.BACKCMD} mkdir -ep {endpoint} -p {dir}" + bcmd = self.valid_command(self.BACKCMD, logact) + cmd = f"{bcmd} mkdir -ep {endpoint} -p {dir}" for loop in range(2): buf = self.pgsystem(cmd, logact, self.CMDRET) syserr = self.PGLOG['SYSERR'] @@ -1408,8 +1419,11 @@ def check_object_file(self, file, bucket = None, opt = 0, logact = 0): if not bucket: bucket = self.PGLOG['OBJCTBKT'] ret = None if not file: return ret - cmd = "{} lo {} -b {}".format(self.OBJCTCMD, file, bucket) - ucmd = "{} gm -k {} -b {}".format(self.OBJCTCMD, file, bucket) if opt&14 else None + ms = re.match(r'^(.+)/$', file) + if ms: file = ms.group(1) # remove ending '/' in case + ocmd = self.valid_command(self.OBJCTCMD, logact) + cmd = "{} lo {} -b {}".format(ocmd, file, bucket) + ucmd = "{} gm -k {} -b {}".format(ocmd, file, bucket) if opt&14 else None loop = 0 while loop < 2: buf = self.pgsystem(cmd, self.LOGWRN, self.CMDRET) @@ -1417,14 +1431,23 @@ def check_object_file(self, file, bucket = None, opt = 0, logact = 0): if re.match(r'^\[\]', buf): break if re.match(r'^\[\{', buf): ary = json.loads(buf) - cnt = len(ary) - if cnt > 1: return self.pglog("{}-{}: {} records returned\n{}".format(bucket, file, cnt, buf), logact|self.ERRLOG) hash = ary[0] uhash = None if ucmd: ubuf = self.pgsystem(ucmd, self.LOGWRN, self.CMDRET) if ubuf and re.match(r'^\{', ubuf): uhash = json.loads(ubuf) ret = self.object_file_stat(hash, uhash, opt) + if ret: + cnt = len(ary) + if cnt > 1 or hash['Key'] != file: + ret['count'] = cnt + ret['fname'] = op.basename(file) + ret['isfile'] = 0 + size = 0 + for a in ary: + size += int(a['Size']) + ret['data_size'] = size + uhash = None break if opt&64: return self.FAILURE errmsg = "Error Execute: {}\n{}".format(cmd, self.PGLOG['SYSERR']) @@ -1443,7 +1466,8 @@ def check_object_path(self, path, bucket = None, logact = 0): if not bucket: bucket = self.PGLOG['OBJCTBKT'] ret = None if not path: return ret - cmd = "{} lo {} -ls -b {}".format(self.OBJCTCMD, path, bucket) + ocmd = self.valid_command(self.OBJCTCMD, logact) + cmd = "{} lo {} -ls -b {}".format(ocmd, path, bucket) loop = 0 while loop < 2: buf = self.pgsystem(cmd, self.LOGWRN, self.CMDRET) @@ -1496,7 +1520,8 @@ def check_backup_file(self, file, endpoint = None, opt = 0, logact = 0): if not endpoint: endpoint = self.PGLOG['BACKUPEP'] bdir = op.dirname(file) bfile = op.basename(file) - cmd = f"{self.BACKCMD} ls -ep {endpoint} -p {bdir} --filter {bfile}" + bcmd = self.valid_command(self.BACKCMD, logact) + cmd = f"{bcmd} ls -ep {endpoint} -p {bdir} --filter {bfile}" ccnt = loop = 0 while loop < 2: buf = self.pgsystem(cmd, logact, self.CMDRET) @@ -1756,7 +1781,8 @@ def object_glob(self, dir, bucket = None, opt = 0, logact = 0): if not bucket: bucket = self.PGLOG['OBJCTBKT'] ms = re.match(r'^(.+)/$', dir) if ms: dir = ms.group(1) - cmd = "{} lo {} -b {}".format(self.OBJCTCMD, dir, bucket) + ocmd = self.valid_command(self.OBJCTCMD, logact) + cmd = "{} lo {} -b {}".format(ocmd, dir, bucket) ary = err = None buf = self.pgsystem(cmd, self.LOGWRN, self.CMDRET) if buf: @@ -1775,7 +1801,7 @@ def object_glob(self, dir, bucket = None, opt = 0, logact = 0): for hash in ary: uhash = None if opt&10: - ucmd = "{} gm -l {} -b {}".format(self.OBJCTCMD, hash['Key'], bucket) + ucmd = "{} gm -l {} -b {}".format(ocmd, hash['Key'], bucket) ubuf = self.pgsystem(ucmd, self.LOGWRN, self.CMDRET) if ubuf and re.match(r'^\{.+', ubuf): uhash = json.loads(ubuf) info = self.object_file_stat(hash, uhash, opt) @@ -1794,7 +1820,8 @@ def object_glob(self, dir, bucket = None, opt = 0, logact = 0): def backup_glob(self, dir, endpoint = None, opt = 0, logact = 0): if not dir: return None if not endpoint: endpoint = self.PGLOG['BACKUPEP'] - cmd = f"{self.BACKCMD} ls -ep {endpoint} -p {dir}" + bcmd = self.valid_command(self.BACKCMD, logact) + cmd = f"{bcmd} ls -ep {endpoint} -p {dir}" flist = {} for loop in range(2): buf = self.pgsystem(cmd, logact, self.CMDRET) diff --git a/src/rda_python_common/pg_log.py b/src/rda_python_common/pg_log.py index fa111fe..26a3821 100644 --- a/src/rda_python_common/pg_log.py +++ b/src/rda_python_common/pg_log.py @@ -101,6 +101,7 @@ def __init__(self): 'OBJCTSTR': "object", 'BACKUPNM': "quasar", 'DRDATANM': "drdata", + 'TACCNAME': "tacc", 'GPFSNAME': "glade", 'PBSNAME': "PBS", 'DSIDCHRS': "d", @@ -866,7 +867,6 @@ def valid_command(self, cmd, logact = 0): # add carbon copies to self.PGLOG['CCDADDR'] def add_carbon_copy(self, cc = None, isstr = None, exclude = 0, specialist = None): - if not cc: if cc is None and isstr is None: self.PGLOG['CCDADDR'] = '' else: @@ -885,28 +885,24 @@ def add_carbon_copy(self, cc = None, isstr = None, exclude = 0, specialist = Non # get the current host name; or batch sever name if getbatch is 1 def get_host(self, getbatch = 0): - if getbatch and self.PGLOG['CURBID'] != 0: host = self.PGLOG['PGBATCH'] elif self.PGLOG['HOSTNAME']: return self.PGLOG['HOSTNAME'] else: host = socket.gethostname() - return self.get_short_host(host) # # strip domain names and retrun the server name itself # def get_short_host(self, host): - if not host: return '' ms = re.match(r'^([^\.]+)\.', host) if ms: host = ms.group(1) if self.PGLOG['HOSTNAME'] and (host == 'localhost' or host == self.PGLOG['HOSTNAME']): return self.PGLOG['HOSTNAME'] HOST = host.upper() if HOST in self.BCHCMDS: return HOST - return host # get a live PBS host name @@ -1063,6 +1059,7 @@ def set_common_pglog(self): self.SETPGLOG("OBJCTBKT", "gdex-data") # default Bucket on Object Store self.SETPGLOG("BACKUPEP", "gdex-quasar") # default Globus Endpoint on Quasar self.SETPGLOG("DRDATAEP", "gdex-quasar-drdata") # DRDATA Globus Endpoint on Quasar + self.SETPGLOG("TACCEP", "gdex-tacc") # default Globus Endpoint on TACC self.SETPGLOG("DBGFILE", "pgdss.dbg") # debug file name self.SETPGLOG("CNFPATH", self.PGLOG['DSSHOME']+"/config") # path to configuration files self.SETPGLOG("DSSURL", "https://gdex.ucar.edu") # current dss web URL