Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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" },
]
Expand Down
81 changes: 54 additions & 27 deletions src/rda_python_common/pg_file.py
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,8 @@ def __init__(self):
'G': self.PGLOG['GPFSNAME'],
'O': self.OHOST,
'B': self.BHOST,
'D': self.DHOST
'D': self.DHOST,

Copilot AI Mar 2, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

DHOSTS adds D but also introduces THOST above; however DHOSTS doesn’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 to self.THOST (and remove the blank line). Otherwise, drop THOST/related constants to avoid partially-wired configuration.

Suggested change
'T': self.THOST,

Copilot uses AI. Check for mistakes.
}
self.DPATHS = {
'G': self.PGLOG['DSSDATA'],
Expand All @@ -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'
}
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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 -'
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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'])
Expand Down Expand Up @@ -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:
Expand All @@ -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)
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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):
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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']
Expand Down Expand Up @@ -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

Copilot AI Mar 2, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

uhash can be referenced before assignment. When ucmd is false (or when ucmd is true but ubuf is empty/non-JSON), uhash is never initialized before being passed to object_file_stat, which will raise UnboundLocalError. Initialize uhash = None before the if ucmd: block (and keep it scoped per-loop).

Copilot uses AI. Check for mistakes.
if ret:
cnt = len(ary)
if cnt > 1 or hash['Key'] != file:
ret['count'] = cnt
Comment on lines 1438 to +1443

Copilot AI Mar 3, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

When check_object_file() receives multiple keys (directory/prefix behavior), ret is first populated via object_file_stat() using only ary[0]. That can leave fields like date_modified, checksum, and meta reflecting the first object even though you later force isfile = 0 and aggregate sizes. Consider clearing those file-specific fields for directory results, or computing directory-appropriate values (e.g., latest mtime across keys) so callers don’t get misleading metadata.

Copilot uses AI. Check for mistakes.
ret['fname'] = op.basename(file)
ret['isfile'] = 0
Comment on lines +1443 to +1445

Copilot AI Mar 3, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

For directory/prefix results you set ret['fname'] = op.basename(file), but op.basename('some/dir/') is an empty string. If callers may pass trailing slashes, consider normalizing file first (e.g., stripping trailing /) before computing fname/comparisons so the returned info isn’t missing a name.

Copilot uses AI. Check for mistakes.
size = 0
Comment on lines +1442 to +1446

Copilot AI Mar 2, 2026

Copy link

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

When treating an object key as a “directory” (cnt > 1 or hash['Key'] != file), the returned ret dict is still initially populated from the first object via object_file_stat (e.g., date_modified, time_modified, checksum, and metadata). Those values will be incorrect/misleading for a directory-like result. Consider returning a separate structure for directory results (or clearing file-specific fields) so callers don’t accidentally consume stale metadata from the first entry.

Copilot uses AI. Check for mistakes.
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'])
Expand All @@ -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)
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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:
Expand All @@ -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)
Expand All @@ -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)
Expand Down
7 changes: 2 additions & 5 deletions src/rda_python_common/pg_log.py
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,7 @@ def __init__(self):
'OBJCTSTR': "object",
'BACKUPNM': "quasar",
'DRDATANM': "drdata",
'TACCNAME': "tacc",
'GPFSNAME': "glade",
'PBSNAME': "PBS",
'DSIDCHRS': "d",
Expand Down Expand Up @@ -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:
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down