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
5 changes: 1 addition & 4 deletions pandaserver/daemons/scripts/copyArchive.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"],
Expand Down Expand Up @@ -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:
Expand All @@ -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"],
Expand Down Expand Up @@ -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"],
Expand Down
2 changes: 0 additions & 2 deletions pandaserver/taskbuffer/TaskBuffer.py
Original file line number Diff line number Diff line change
Expand Up @@ -558,7 +558,6 @@ def storeJobs(
# lock jobs for reassign
def lockJobsForReassign(
self,
tableName: str,
timeLimit: datetime.datetime,
statList: list[str],
labels: list[str],
Expand All @@ -573,7 +572,6 @@ def lockJobsForReassign(
with self.proxyPool.get() as proxy:
# exec
res = proxy.lockJobsForReassign(
tableName,
timeLimit,
statList,
labels,
Expand Down
19 changes: 13 additions & 6 deletions pandaserver/taskbuffer/db_proxy_mods/job_standalone_module.py
Original file line number Diff line number Diff line change
Expand Up @@ -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],
Expand All @@ -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} "
Expand All @@ -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}) "
Expand Down
Loading