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_dsquasar"
version = "3.0.15"
version = "3.0.16"
authors = [
{ name="Zaihua Ji", email="zji@ucar.edu" },
]
Expand Down
26 changes: 23 additions & 3 deletions src/rda_python_dsquasar/dsquasar.py
Original file line number Diff line number Diff line change
Expand Up @@ -611,7 +611,7 @@ def backup_dataset_files(self, dsfiles, backflag):
qinfo['dslocks'].remove(dsid)
qinfo['dsids'] = []
if self.PGBACK['mproc'] > 1:
self.check_child(None, 0, self.LOGWRN, 1) # wait all child processes done
self.wait_all_children() # wait all child processes done
self.confirm_quasar_counts(qinfo, qinfo['abids'], "status = 'T'") # recount confirmed backups from RDADB
if qinfo['dslocks']:
for dsid in qinfo['dslocks']:
Expand Down Expand Up @@ -728,7 +728,7 @@ def backup_dataset_infiles(self, dsfiles, backflag):
self.unlock_backup_datasets(qinfo, dslocks)
if self.PGBACK['dolock']: self.unlock_backup_datasets(qinfo, [dsid])
if self.PGBACK['mproc'] > 1:
self.check_child(None, 0, self.LOGWRN, 1) # wait all child processes done
self.wait_all_children() # wait all child processes done
self.confirm_quasar_counts(qinfo, qinfo['abids'], "status = 'T'") # recount confirmed backups from RDADB
qcnt = qinfo['qcnt']
if qcnt > 0:
Expand Down Expand Up @@ -769,7 +769,7 @@ def backup_dataset_tarfiles(self, dsfiles, backflag):
self.transfer_quasar_tarfiles(qinfo)
if qinfo['size'] > 0: self.transfer_quasar_tarfiles(qinfo)
if self.PGBACK['mproc'] > 1:
self.check_child(None, 0, self.LOGWRN, 1) # wait all child processes done
self.wait_all_children() # wait all child processes done
self.confirm_quasar_counts(qinfo, list(bfiles), "status = 'A'") # recount confirmed backups from RDADB
qcnt = qinfo['qcnt']
if qcnt > 0:
Expand Down Expand Up @@ -906,6 +906,20 @@ def check_batch_deadline(self):
self.report_dscheck_email(title, 1)
self.PGLOG.update(saved)

# wait for every child process to finish, checking the walltime deadline between polls.
# check_child(..., 1) does the waiting inside its own loop, so a phase that ends with the
# last tar children still running would sit there for the rest of the run without ever
# reaching a deadline check - the parent has no per-item loop left to guard. polling with
# dowait 0 and sleeping here keeps the same cadence and wait message as check_child does.
def wait_all_children(self):
if self.PGBACK['mproc'] < 2: return
i = 0
while True:
self.check_batch_deadline()
if self.check_child(None, 0, self.LOGWRN, 0) < 1: break
self.show_wait_message(i, "{}: wait child processes".format(self.PGSIG['DSTR']), self.LOGWRN, 1)
i += 1

# recompute the confirmed backup counts from RDADB after all child processes
# finished, so a multi-process summary reflects succeeded (not just started) work
def confirm_quasar_counts(self, qinfo, bids, dcnd):
Expand Down Expand Up @@ -1389,6 +1403,10 @@ def gather_dataset_files(self, dsfiles, unlock = True, sizes = None):
dcnt = len(mrecs['dsid']) if mrecs else 0
pgrecs = [self.onerecord(mrecs, i) for i in range(dcnt)]
for pgrec in pgrecs:
# a gather is not free: it queries every dataset, and each one another worker holds
# costs an rdaps call to check that worker's pid. a run under dscheck does two of
# them before any work, and -A 3 does a third one before tarring
self.check_batch_deadline()
dsid = pgrec['dsid']
if unlock and pgrec['pid'] and self.lock_dataset(dsid, 0, self.LOGACT) < 1: continue
fcnt += self.get_dataset_files(dsid, dsfiles, pgrec['backflag'], logact, sizes)
Expand All @@ -1413,6 +1431,7 @@ def gather_dataset_infiles(self, dsfiles):
pgrecs = self.pgmget("bfile", flds, dcnd, self.LGWNEX)
bcnt = len(pgrecs['bid']) if pgrecs else 0
for i in range(bcnt):
self.check_batch_deadline()
cnt = self.get_backup_infile(dsfiles, self.onerecord(pgrecs, i))
if cnt:
icnt += 1
Expand All @@ -1438,6 +1457,7 @@ def gather_dataset_tarfiles(self, dsfiles):
pgrecs = self.pgmget("bfile", flds, dcnd, self.LGWNEX)
bcnt = len(pgrecs['bid']) if pgrecs else 0
for i in range(bcnt):
self.check_batch_deadline()
cnt = self.get_backup_tarfile(dsfiles, self.onerecord(pgrecs, i))
if cnt:
tcnt += 1
Expand Down
Loading