From bf2812ba47eb133743e7013bf619a0280a4141b8 Mon Sep 17 00:00:00 2001 From: tmaeno Date: Thu, 24 Sep 2026 14:50:53 +0200 Subject: [PATCH] Refactor lockJobsForReassign method to resolve table names based on job statuses --- pandaserver/daemons/scripts/copyArchive.py | 5 +---- pandaserver/taskbuffer/TaskBuffer.py | 2 -- .../db_proxy_mods/job_standalone_module.py | 19 +++++++++++++------ 3 files changed, 14 insertions(+), 12 deletions(-) diff --git a/pandaserver/daemons/scripts/copyArchive.py b/pandaserver/daemons/scripts/copyArchive.py index 05fbb9884..5cea40624 100644 --- a/pandaserver/daemons/scripts/copyArchive.py +++ b/pandaserver/daemons/scripts/copyArchive.py @@ -725,7 +725,6 @@ def _memoryCheck(str: str) -> None: timeLimit = naive_utcnow() - datetime.timedelta(minutes=timeoutValue) # get PandaIDs status, res = taskBuffer.lockJobsForReassign( - "ATLAS_PANDA.jobsDefined4", timeLimit, ["defined"], ["managed", "test"], @@ -772,7 +771,7 @@ def _memoryCheck(str: str) -> None: # reassign long-waiting jobs in defined table timeLimit = naive_utcnow() - datetime.timedelta(hours=12) - status, res = taskBuffer.lockJobsForReassign("ATLAS_PANDA.jobsDefined4", timeLimit, [], ["managed"], [], [], [], True) + status, res = taskBuffer.lockJobsForReassign(timeLimit, ["defined", "assigned", "waiting", "pending"], ["managed"], [], [], [], True) jediJobs = [] if res is not None: for id, lockedby in res: @@ -793,7 +792,6 @@ def _memoryCheck(str: str) -> None: # reassign too long activated jobs in active table timeLimit = naive_utcnow() - datetime.timedelta(days=2) status, res = taskBuffer.lockJobsForReassign( - "ATLAS_PANDA.jobsActive4", timeLimit, ["activated"], ["managed"], @@ -834,7 +832,6 @@ def _memoryCheck(str: str) -> None: # reassign too long starting jobs in active table timeLimit = naive_utcnow() - datetime.timedelta(hours=48) status, res = taskBuffer.lockJobsForReassign( - "ATLAS_PANDA.jobsActive4", timeLimit, ["starting"], ["managed"], diff --git a/pandaserver/taskbuffer/TaskBuffer.py b/pandaserver/taskbuffer/TaskBuffer.py index 8f2460e9a..fbaf89fbe 100755 --- a/pandaserver/taskbuffer/TaskBuffer.py +++ b/pandaserver/taskbuffer/TaskBuffer.py @@ -558,7 +558,6 @@ def storeJobs( # lock jobs for reassign def lockJobsForReassign( self, - tableName: str, timeLimit: datetime.datetime, statList: list[str], labels: list[str], @@ -573,7 +572,6 @@ def lockJobsForReassign( with self.proxyPool.get() as proxy: # exec res = proxy.lockJobsForReassign( - tableName, timeLimit, statList, labels, diff --git a/pandaserver/taskbuffer/db_proxy_mods/job_standalone_module.py b/pandaserver/taskbuffer/db_proxy_mods/job_standalone_module.py index a606114e3..82794c66c 100644 --- a/pandaserver/taskbuffer/db_proxy_mods/job_standalone_module.py +++ b/pandaserver/taskbuffer/db_proxy_mods/job_standalone_module.py @@ -774,7 +774,6 @@ def setDebugMode(self, dn: str, pandaID: int, prodManager: bool, modeOn: bool, w # lock jobs for reassign def lockJobsForReassign( self, - tableName: str, timeLimit: datetime.datetime, statList: list[str], labels: list[str], @@ -788,8 +787,17 @@ def lockJobsForReassign( ) -> tuple[bool, list[Any]]: comment = " /* DBProxy.lockJobsForReassign */" tmp_log = self.create_tagged_logger(comment) - tmp_log.debug(f"{tableName} {timeLimit} {statList} {labels} {processTypes} {sites} {clouds} {useJEDI}") + tmp_log.debug(f"{timeLimit} {statList} {labels} {processTypes} {sites} {clouds} {useJEDI}") try: + # resolve table from job statuses + defined_statuses = {"defined", "assigned", "waiting", "pending"} + active_statuses = {"activated", "throttled", "sent", "starting", "running", "holding", "transferring", "merging"} + if statList and set(statList) <= defined_statuses: + tableName = "ATLAS_PANDA.jobsDefined4" + elif statList and set(statList) <= active_statuses: + tableName = "ATLAS_PANDA.jobsActive4" + else: + raise ValueError(f"cannot resolve table from statList={statList}") # make sql if not useJEDI: sql = f"SELECT PandaID FROM {tableName} " @@ -803,10 +811,9 @@ def lockJobsForReassign( sql += "WHERE stateChangeTime<:modificationTime " varMap: dict[str, Any] = {} varMap[":modificationTime"] = timeLimit - if statList != []: - stat_var_names_str, stat_var_map = get_sql_IN_bind_variables(statList, prefix=":stat") - sql += f"AND jobStatus IN ({stat_var_names_str}) " - varMap.update(stat_var_map) + stat_var_names_str, stat_var_map = get_sql_IN_bind_variables(statList, prefix=":stat") + sql += f"AND jobStatus IN ({stat_var_names_str}) " + varMap.update(stat_var_map) if labels != []: label_var_names_str, label_var_map = get_sql_IN_bind_variables(labels, prefix=":label") sql += f"AND prodSourceLabel IN ({label_var_names_str}) "