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 README.md
Original file line number Diff line number Diff line change
Expand Up @@ -165,7 +165,7 @@ PgLOG.pglog("hello", PgLOG.LOGWRN)
python -c "import rda_python_common; print(rda_python_common.__version__)"
```

You should see the installed version (currently `3.0.18`). If the import
You should see the installed version (currently `3.0.19`). If the import
fails, double-check that the active Python environment is the one where you
ran `pip install`.

Expand Down
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 = "3.0.18"
version = "3.0.19"
authors = [
{ name="Zaihua Ji", email="zji@ucar.edu" },
]
Expand Down
2 changes: 1 addition & 1 deletion src/rda_python_common/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@

from . import PgLOG, PgUtil, PgDBI, PgFile, PgLock, PgCMD, PgSIG, PgOPT, PgSplit

__version__ = "3.0.18"
__version__ = "3.0.19"

__all__ = [
"PgLOG",
Expand Down
8 changes: 8 additions & 0 deletions src/rda_python_common/pg_cmd.py
Original file line number Diff line number Diff line change
Expand Up @@ -207,6 +207,14 @@ def init_dscheck(self, oindex, otype, cmd, dsid, action, workdir=None, specialis
self.pglog(cmsg + "is Running, No restart", self.LOGWRN)
sys.exit(0)
if cidx > 0:
# a report parked in einfo by the previous run is only mailed by the dscheck
# daemon, which skips any record that is already running. restarting the record
# here would strand that report until the next run overwrites it, so leave the
# record alone and let the daemon mail it first - it clears einfo when it does
if pgrec['einfo']:
self.pglog(cmsg + "has a report to email, No restart", self.LOGWRN)
self.lock_dscheck(cidx, 0, logact)
sys.exit(0)
if not hosts and pgrec['hostname']:
hosts = pgrec['hostname']
self.set_one_boption('hostname', hosts, 0)
Expand Down
26 changes: 21 additions & 5 deletions src/rda_python_common/pg_file.py
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,8 @@ def __init__(self):
}
self.TARSTR = '|'.join(self.PGTARS)
self.DELDIRS = {}
self.TASKIDS = {} # cache unfinished
self.TASKIDS = {} # cache unfinished
self.QCANCEL = {} # taskid -> why a task got cancelled, reported by the waiting caller
self.LHOST = "localhost"
self.OHOST = self.PGLOG['OBJCTSTR']
self.BHOST = self.PGLOG['BACKUPNM']
Expand Down Expand Up @@ -532,7 +533,11 @@ def submit_globus_task(self, cmd, endpoint, logact = 0, qstr = None):
time.sleep(self.PGSIG['ETIME'])
lp += 1
if task['stat'] == 'S' or task['stat'] == 'A': break
if task['stat'] == 'F' and not syserr: break
if task['stat'] == 'F' and not syserr:
# nothing waits on this task, so report the cancellation here instead
if task['id'] in self.QCANCEL:
self.errlog("{}: Cancel Task due to {}".format(task['id'], self.QCANCEL.pop(task['id'])), 'B', 1, logact)
break
Comment on lines +536 to +540
errmsg = "Error Execute: " + cmd
if qstr: errmsg += " with stdin:\n" + qstr
if syserr:
Expand Down Expand Up @@ -581,8 +586,16 @@ def check_globus_status(self, taskid, endpoint = None, logact = 0):
detail = ms.group(1)
if detail not in astats:
if logact&self.NOWAIT:
errmsg = "{}: Cancel Task due to {}:\n{}".format(taskid, detail, buf)
self.errlog(errmsg, 'B', 1, logact)
# record why, and let the caller waiting on this task report it as a
# single error naming the file; dumping the whole get-task output here
# doubled every failure into two error entries, the first 17 lines long
reason = detail
ms = re.search(r'Bytes Transferred:\s+(\d+)', buf)
if ms: reason += " after " + self.format_float_value(ms.group(1))
ms = re.search(r'Files:\s+(\d+)', buf)
if ms: reason += " of {} file(s)".format(ms.group(1))
self.QCANCEL[taskid] = reason
self.pglog("{}: Cancel Task due to {}".format(taskid, reason), self.LOGWRN)
ccmd = f"{bcmd} cancel-task {taskid}"
self.pgsystem(ccmd, logact, 7)
else:
Expand Down Expand Up @@ -640,7 +653,10 @@ def check_globus_finished(self, tofile, topoint, logact = 0):
del self.TASKIDS[ckey]
else:
status = self.QSTATS[stat] if stat in self.QSTATS else 'UNKNOWN'
self.errlog("{}: Status '{}' for Task {}".format(ckey, status, taskid), 'B', 1, logact)
errmsg = "{}: Status '{}' for Task {}".format(ckey, status, taskid)
if taskid in self.QCANCEL:
errmsg += " - " + self.QCANCEL.pop(taskid)
self.errlog(errmsg, 'B', 1, logact)
ret = self.FAILURE
break
return ret
Expand Down
Loading