-
Notifications
You must be signed in to change notification settings - Fork 0
check object for directory #85
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -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,15 +253,16 @@ 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) | ||
| if not logact&self.OVRIDE: | ||
| 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,23 +1419,35 @@ 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) | ||
| if buf: | ||
| 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) | ||
|
Comment on lines
1434
to
1439
|
||
| if ret: | ||
| cnt = len(ary) | ||
| if cnt > 1 or hash['Key'] != file: | ||
| ret['count'] = cnt | ||
|
Comment on lines
1438
to
+1443
|
||
| ret['fname'] = op.basename(file) | ||
| ret['isfile'] = 0 | ||
|
Comment on lines
+1443
to
+1445
|
||
| size = 0 | ||
|
Comment on lines
+1442
to
+1446
|
||
| 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) | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
DHOSTSaddsDbut also introducesTHOSTabove; howeverDHOSTSdoesn’t include an entry for the new TACC host (and contains a stray blank line). If TACC is intended to be a supported down-storage flag, add the corresponding key (e.g.,T) mapped toself.THOST(and remove the blank line). Otherwise, dropTHOST/related constants to avoid partially-wired configuration.