Coverage for python/lsst/ctrl/bps/htcondor/lssthtc.py: 79%

954 statements  

« prev     ^ index     » next       coverage.py v7.16.1, created at 2026-09-25 22:30 +0000

1# This file is part of ctrl_bps_htcondor. 

2# 

3# Developed for the LSST Data Management System. 

4# This product includes software developed by the LSST Project 

5# (https://www.lsst.org). 

6# See the COPYRIGHT file at the top-level directory of this distribution 

7# for details of code ownership. 

8# 

9# This software is dual licensed under the GNU General Public License and also 

10# under a 3-clause BSD license. Recipients may choose which of these licenses 

11# to use; please see the files gpl-3.0.txt and/or bsd_license.txt, 

12# respectively. If you choose the GPL option then the following text applies 

13# (but note that there is still no warranty even if you opt for BSD instead): 

14# 

15# This program is free software: you can redistribute it and/or modify 

16# it under the terms of the GNU General Public License as published by 

17# the Free Software Foundation, either version 3 of the License, or 

18# (at your option) any later version. 

19# 

20# This program is distributed in the hope that it will be useful, 

21# but WITHOUT ANY WARRANTY; without even the implied warranty of 

22# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the 

23# GNU General Public License for more details. 

24# 

25# You should have received a copy of the GNU General Public License 

26# along with this program. If not, see <https://www.gnu.org/licenses/>. 

27 

28"""Placeholder HTCondor DAGMan API. 

29 

30There is new work on a python DAGMan API from HTCondor. However, at this 

31time, it tries to make things easier by assuming DAG is easily broken into 

32levels where there are 1-1 or all-to-all relationships to nodes in next 

33level. LSST workflows are more complicated. 

34""" 

35 

36__all__ = [ 

37 "MISSING_ID", 

38 "DagStatus", 

39 "HTCDag", 

40 "HTCJob", 

41 "NodeStatus", 

42 "RestrictedDict", 

43 "WmsNodeType", 

44 "condor_history", 

45 "condor_q", 

46 "condor_search", 

47 "condor_status", 

48 "htc_backup_files", 

49 "htc_check_dagman_output", 

50 "htc_create_submit_from_cmd", 

51 "htc_create_submit_from_dag", 

52 "htc_create_submit_from_file", 

53 "htc_escape", 

54 "htc_query_history", 

55 "htc_query_present", 

56 "htc_submit_dag", 

57 "htc_tweak_log_info", 

58 "htc_version", 

59 "htc_write_attribs", 

60 "htc_write_condor_file", 

61 "pegasus_name_to_label", 

62 "read_dag_info", 

63 "read_dag_log", 

64 "read_dag_nodes_log", 

65 "read_dag_status", 

66 "read_node_status", 

67 "summarize_dag", 

68 "update_job_info", 

69 "write_dag_info", 

70] 

71 

72 

73import itertools 

74import json 

75import logging 

76import os 

77import pprint 

78import re 

79import subprocess 

80from collections import Counter, defaultdict 

81from collections.abc import MutableMapping 

82from datetime import datetime, timedelta 

83from enum import IntEnum, auto 

84from pathlib import Path 

85from typing import Any, TextIO 

86 

87import classad 

88import htcondor 

89import networkx 

90from deprecated.sphinx import deprecated 

91from packaging import version 

92 

93from .handlers import HTC_JOB_AD_HANDLERS 

94 

95_LOG = logging.getLogger(__name__) 

96 

97MISSING_ID = "-99999" 

98 

99 

100class DagStatus(IntEnum): 

101 """HTCondor DAGMan's statuses for a DAG.""" 

102 

103 OK = 0 

104 ERROR = 1 # an error condition different than those listed here 

105 FAILED = 2 # one or more nodes in the DAG have failed 

106 ABORTED = 3 # the DAG has been aborted by an ABORT-DAG-ON specification 

107 REMOVED = 4 # the DAG has been removed by condor_rm 

108 CYCLE = 5 # a cycle was found in the DAG 

109 SUSPENDED = 6 # the DAG has been suspended (see section 2.10.8) 

110 

111 

112@deprecated( 

113 reason="The JobStatus is internally replaced by htcondor.JobStatus. " 

114 "External reporting code should be using ctrl_bps.WmsStates. " 

115 "This class will be removed after v30.", 

116 version="v30.0", 

117 category=FutureWarning, 

118) 

119class JobStatus(IntEnum): 

120 """HTCondor's statuses for jobs.""" 

121 

122 UNEXPANDED = 0 # Unexpanded 

123 IDLE = 1 # Idle 

124 RUNNING = 2 # Running 

125 REMOVED = 3 # Removed 

126 COMPLETED = 4 # Completed 

127 HELD = 5 # Held 

128 TRANSFERRING_OUTPUT = 6 # Transferring_Output 

129 SUSPENDED = 7 # Suspended 

130 

131 

132class NodeStatus(IntEnum): 

133 """HTCondor's statuses for DAGman nodes.""" 

134 

135 # (STATUS_NOT_READY): At least one parent has not yet finished or the node 

136 # is a FINAL node. 

137 NOT_READY = 0 

138 

139 # (STATUS_READY): All parents have finished, but the node is not yet 

140 # running. 

141 READY = 1 

142 

143 # (STATUS_PRERUN): The node’s PRE script is running. 

144 PRERUN = 2 

145 

146 # (STATUS_SUBMITTED): The node’s HTCondor job(s) are in the queue. 

147 # StatusDetails = "not_idle" -> running. 

148 # JobProcsHeld = 1-> hold. 

149 # JobProcsQueued = 1 -> idle. 

150 SUBMITTED = 3 

151 

152 # (STATUS_POSTRUN): The node’s POST script is running. 

153 POSTRUN = 4 

154 

155 # (STATUS_DONE): The node has completed successfully. 

156 DONE = 5 

157 

158 # (STATUS_ERROR): The node has failed. StatusDetails has info (e.g., 

159 # ULOG_JOB_ABORTED for deleted job). 

160 ERROR = 6 

161 

162 # (STATUS_FUTILE): The node will never run because ancestor node failed. 

163 FUTILE = 7 

164 

165 

166class WmsNodeType(IntEnum): 

167 """HTCondor plugin node types to help with payload reporting.""" 

168 

169 UNKNOWN = auto() 

170 """Dummy value when missing.""" 

171 

172 PAYLOAD = auto() 

173 """Payload job.""" 

174 

175 FINAL = auto() 

176 """Final job.""" 

177 

178 SERVICE = auto() 

179 """Service job.""" 

180 

181 NOOP = auto() 

182 """NOOP job used for ordering jobs.""" 

183 

184 SUBDAG = auto() 

185 """SUBDAG job used for ordering jobs.""" 

186 

187 SUBDAG_CHECK = auto() 

188 """Job used to correctly prune jobs after a subdag.""" 

189 

190 

191HTC_QUOTE_KEYS = {"environment", "arguments"} 

192HTC_VALID_JOB_KEYS = { 

193 "universe", 

194 "executable", 

195 "arguments", 

196 "environment", 

197 "log", 

198 "error", 

199 "output", 

200 "should_transfer_files", 

201 "when_to_transfer_output", 

202 "getenv", 

203 "notification", 

204 "notify_user", 

205 "concurrency_limit", 

206 "transfer_executable", 

207 "transfer_input_files", 

208 "transfer_output_files", 

209 "transfer_output_remaps", 

210 "request_cpus", 

211 "request_memory", 

212 "request_disk", 

213 "priority", 

214 "category", 

215 "requirements", 

216 "on_exit_hold", 

217 "on_exit_hold_reason", 

218 "on_exit_hold_subcode", 

219 "max_retries", 

220 "retry_until", 

221 "periodic_release", 

222 "periodic_remove", 

223 "accounting_group", 

224 "accounting_group_user", 

225 "kill_sig", 

226 "want_graceful_removal", 

227 "job_max_vacate_time", 

228} 

229HTC_VALID_JOB_DAG_KEYS = { 

230 "dir", 

231 "noop", 

232 "done", 

233 "vars", 

234 "pre", 

235 "post", 

236 "retry", 

237 "retry_unless_exit", 

238 "abort_dag_on", 

239 "abort_exit", 

240 "priority", 

241} 

242HTC_VERSION = version.parse(htcondor.__version__) 

243 

244 

245class RestrictedDict(MutableMapping): 

246 """A dictionary that only allows certain keys. 

247 

248 Parameters 

249 ---------- 

250 valid_keys : `~collections.abc.Container` 

251 Strings that are valid keys. 

252 init_data : `dict` or `RestrictedDict`, optional 

253 Initial data. 

254 

255 Raises 

256 ------ 

257 KeyError 

258 If invalid key(s) in init_data. 

259 """ 

260 

261 def __init__(self, valid_keys, init_data=()): 

262 self.valid_keys = valid_keys 

263 self.data = {} 

264 self.update(init_data) 

265 

266 def __getitem__(self, key): 

267 """Return value for given key if exists. 

268 

269 Parameters 

270 ---------- 

271 key : `str` 

272 Identifier for value to return. 

273 

274 Returns 

275 ------- 

276 value : `~typing.Any` 

277 Value associated with given key. 

278 

279 Raises 

280 ------ 

281 KeyError 

282 If key doesn't exist. 

283 """ 

284 return self.data[key] 

285 

286 def __delitem__(self, key): 

287 """Delete value for given key if exists. 

288 

289 Parameters 

290 ---------- 

291 key : `str` 

292 Identifier for value to delete. 

293 

294 Raises 

295 ------ 

296 KeyError 

297 If key doesn't exist. 

298 """ 

299 del self.data[key] 

300 

301 def __setitem__(self, key, value): 

302 """Store key,value in internal dict only if key is valid. 

303 

304 Parameters 

305 ---------- 

306 key : `str` 

307 Identifier to associate with given value. 

308 value : `~typing.Any` 

309 Value to store. 

310 

311 Raises 

312 ------ 

313 KeyError 

314 If key is invalid. 

315 """ 

316 if key not in self.valid_keys: 316 ↛ 317line 316 didn't jump to line 317 because the condition on line 316 was never true

317 raise KeyError(f"Invalid key {key}") 

318 self.data[key] = value 

319 

320 def __iter__(self): 

321 return self.data.__iter__() 

322 

323 def __len__(self): 

324 return len(self.data) 

325 

326 def __str__(self): 

327 return str(self.data) 

328 

329 

330def htc_backup_files( 

331 wms_path: str | os.PathLike, 

332 subdir: str | os.PathLike | None = None, 

333 limit: int = 100, 

334 failed_subdags: set[str] | None = None, 

335) -> None: 

336 """Backup select HTCondor files in the submit directory. 

337 

338 Files will be saved in separate subdirectories which will be created in 

339 the submit directory where the files are located. These subdirectories 

340 will be consecutive, zero-padded integers. Their values will correspond to 

341 the number of HTCondor rescue DAGs in the submit directory. 

342 

343 Hence, with the default settings, copies after the initial failed run will 

344 be placed in '001' subdirectory, '002' after the first restart, and so on 

345 until the limit of backups is reached. If there's no rescue DAG yet, files 

346 will be copied to '000' subdirectory. 

347 

348 This is not a generic function for making backups. It is intended to be 

349 used once, just before a restart, to make snapshots of files which will be 

350 overwritten by HTCondor after during the next run. 

351 

352 Parameters 

353 ---------- 

354 wms_path : `str` or `os.PathLike` 

355 Path to the submit directory either absolute or relative. 

356 subdir : `str` or `os.PathLike`, optional 

357 A path, relative to the submit directory, where all subdirectories with 

358 backup files will be kept. Defaults to None which means that the backup 

359 subdirectories will be placed directly in the submit directory. 

360 limit : `int`, optional 

361 Maximal number of backups. If the number of backups reaches the limit, 

362 the last backup files will be overwritten. The default value is 100 

363 to match the default value of HTCondor's DAGMAN_MAX_RESCUE_NUM in 

364 version 8.8+. 

365 failed_subdags : `set` [`str`], optional 

366 Names of subdag jobs that failed. Only files of these subdags will be 

367 backed up. If None (the default), files of all subdags with rescue 

368 files are backed up. 

369 

370 Raises 

371 ------ 

372 FileNotFoundError 

373 If the submit directory or the file that needs to be backed up does not 

374 exist. 

375 OSError 

376 If the submit directory cannot be accessed or backing up a file failed 

377 either due to permission or filesystem related issues. 

378 """ 

379 width = len(str(limit)) 

380 

381 path = Path(wms_path).resolve() 

382 if not path.is_dir(): 

383 raise FileNotFoundError(f"Directory {path} not found") 

384 

385 # Initialize the backup counter. 

386 # If using control DAG, don't want to include nested DAGs. 

387 rescue_dags = list(path.glob("*_ctrl.dag.rescue[0-9][0-9][0-9]")) 

388 if not rescue_dags: 388 ↛ 390line 388 didn't jump to line 390 because the condition on line 388 was always true

389 rescue_dags = list(path.glob("*.rescue[0-9][0-9][0-9]")) 

390 counter = min(len(rescue_dags), limit) 

391 

392 # Create the backup directory and move select files there. 

393 dest = path 

394 if subdir: 

395 # PurePath.is_relative_to() is not available before Python 3.9. Hence 

396 # we need to check is 'subdir' is in the submit directory in some other 

397 # way if it is an absolute path. 

398 subdir = Path(subdir) 

399 if subdir.is_absolute(): 

400 subdir = subdir.resolve() # Since resolve was run on path, must run it here 

401 if dest not in subdir.parents: 

402 _LOG.warning( 

403 "Invalid backup location: '%s' not in the submit directory, will use '%s' instead.", 

404 subdir, 

405 wms_path, 

406 ) 

407 else: 

408 dest /= subdir 

409 else: 

410 dest /= subdir 

411 dest /= f"{counter:0{width}}" 

412 _LOG.debug("dest = %s", dest) 

413 try: 

414 dest.mkdir(parents=True, exist_ok=False if counter < limit else True) 

415 except FileExistsError: 

416 _LOG.warning("Refusing to do backups: target directory '%s' already exists", dest) 

417 else: 

418 htc_backup_files_single_path(path, dest) 

419 

420 # Back up selected files for failed subdags as well. 

421 # 

422 # Do NOT back up files of subdags that succeeded! These files need to stay 

423 # in their respective directories as they will not be recreated by HTCondor 

424 # after the run is restarted. As HTCondorService.report() uses information 

425 # in these files to determine job statuses, their absence may lead to 

426 # reporting incorrect job status counts. 

427 for subdag_dir in {file.parent for file in path.glob("subdags/*/*.rescue*")}: 

428 if failed_subdags is not None and subdag_dir.name not in failed_subdags: 428 ↛ 429line 428 didn't jump to line 429 because the condition on line 428 was never true

429 continue 

430 subdag_dest = dest / subdag_dir.relative_to(path) 

431 subdag_dest.mkdir(parents=True, exist_ok=False) 

432 htc_backup_files_single_path(subdag_dir, subdag_dest) 

433 

434 

435def htc_backup_files_single_path(src: str | os.PathLike, dest: str | os.PathLike) -> None: 

436 """Move particular htc files to a different directory for later debugging. 

437 

438 Parameters 

439 ---------- 

440 src : `str` or `os.PathLike` 

441 Directory from which to back up particular files. 

442 dest : `str` or `os.PathLike` 

443 Directory to which particular files are moved. 

444 

445 Raises 

446 ------ 

447 RuntimeError 

448 If given dest directory matches given src directory. 

449 OSError 

450 If problems moving file. 

451 FileNotFoundError 

452 Item matching pattern in src directory isn't a file. 

453 """ 

454 src = Path(src) 

455 dest = Path(dest) 

456 if dest.samefile(src): 

457 raise RuntimeError(f"Destination directory is same as the source directory ({src})") 

458 

459 for patt in [ 

460 "*.info.*", 

461 "*.dag.metrics", 

462 "*.dag.nodes.log", 

463 "*.node_status", 

464 "wms_*.dag.post.out", 

465 "wms_*.status.txt", 

466 ]: 

467 for source in src.glob(patt): 

468 if source.is_file(): 468 ↛ 475line 468 didn't jump to line 475 because the condition on line 468 was always true

469 target = dest / source.relative_to(src) 

470 try: 

471 source.rename(target) 

472 except OSError as exc: 

473 raise type(exc)(f"Backing up '{source}' failed: {exc.strerror}") from None 

474 else: 

475 raise FileNotFoundError(f"Backing up '{source}' failed: not a file") 

476 

477 

478def htc_escape(value): 

479 """Escape characters in given value based upon HTCondor syntax. 

480 

481 Parameters 

482 ---------- 

483 value : `~typing.Any` 

484 Value that needs to have characters escaped if string. 

485 

486 Returns 

487 ------- 

488 new_value : `~typing.Any` 

489 Given value with characters escaped appropriate for HTCondor if string. 

490 """ 

491 if isinstance(value, str): 

492 newval = value.replace('"', '""').replace("'", "''").replace("&quot;", '"') 

493 else: 

494 newval = value 

495 

496 return newval 

497 

498 

499def htc_write_attribs(stream, attrs): 

500 """Write job attributes in HTCondor format to writeable stream. 

501 

502 Parameters 

503 ---------- 

504 stream : `~typing.TextIO` 

505 Output text stream (typically an open file). 

506 attrs : `dict` 

507 HTCondor job attributes (dictionary of attribute key, value). 

508 """ 

509 for key, value in attrs.items(): 

510 # Make sure strings are syntactically correct for HTCondor. 

511 if isinstance(value, str): 511 ↛ 514line 511 didn't jump to line 514 because the condition on line 511 was always true

512 pval = f'"{htc_escape(value)}"' 

513 else: 

514 pval = value 

515 

516 print(f"+{key} = {pval}", file=stream) 

517 

518 

519def htc_write_condor_file( 

520 filename: str | os.PathLike, job_name: str, job: RestrictedDict, job_attrs: dict[str, Any] 

521) -> None: 

522 """Write an HTCondor submit file. 

523 

524 Parameters 

525 ---------- 

526 filename : `str` or `os.PathLike` 

527 Filename for the HTCondor submit file. 

528 job_name : `str` 

529 Job name to use in submit file. 

530 job : `RestrictedDict` 

531 Submit script information. 

532 job_attrs : `dict` 

533 Job attributes. 

534 """ 

535 os.makedirs(os.path.dirname(filename), exist_ok=True) 

536 with open(filename, "w") as fh: 

537 for key, value in job.items(): 

538 if value is not None: 538 ↛ 537line 538 didn't jump to line 537 because the condition on line 538 was always true

539 if key in HTC_QUOTE_KEYS: # Assumes internal quotes are already escaped correctly 

540 print(f'{key}="{value}"', file=fh) 

541 else: 

542 print(f"{key}={value}", file=fh) 

543 for key in ["output", "error", "log"]: 

544 if key not in job: 

545 filename = f"{job_name}.$(Cluster).{'out' if key != 'log' else key}" 

546 print(f"{key}={filename}", file=fh) 

547 

548 if job_attrs is not None: 

549 htc_write_attribs(fh, job_attrs) 

550 print("queue", file=fh) 

551 

552 

553# To avoid doing the version check during every function call select 

554# appropriate conversion function at the import time. 

555# 

556# Make sure that *each* version specific variant of the conversion function(s) 

557# has the same signature after applying any changes! 

558if HTC_VERSION < version.parse("8.9.8"): 558 ↛ 560line 558 didn't jump to line 560 because the condition on line 558 was never true

559 

560 def htc_tune_schedd_args(**kwargs): 

561 """Ensure that arguments for Schedd are version appropriate. 

562 

563 The old arguments: 'requirements' and 'attr_list' of 

564 'Schedd.history()', 'Schedd.query()', and 'Schedd.xquery()' were 

565 deprecated in favor of 'constraint' and 'projection', respectively, 

566 starting from version 8.9.8. The function will convert "new" keyword 

567 arguments to "old" ones. 

568 

569 Parameters 

570 ---------- 

571 **kwargs 

572 Any keyword arguments that Schedd.history(), Schedd.query(), and 

573 Schedd.xquery() accepts. 

574 

575 Returns 

576 ------- 

577 kwargs : `dict` [`str`, `~typing.Any`] 

578 Keywords arguments that are guaranteed to work with the Python 

579 HTCondor API. 

580 

581 Notes 

582 ----- 

583 Function doesn't validate provided keyword arguments beyond converting 

584 selected arguments to their version specific form. For example, 

585 it won't remove keywords that are not supported by the methods 

586 mentioned earlier. 

587 """ 

588 translation_table = { 

589 "constraint": "requirements", 

590 "projection": "attr_list", 

591 } 

592 for new, old in translation_table.items(): 

593 try: 

594 kwargs[old] = kwargs.pop(new) 

595 except KeyError: 

596 pass 

597 return kwargs 

598 

599else: 

600 

601 def htc_tune_schedd_args(**kwargs): 

602 """Ensure that arguments for Schedd are version appropriate. 

603 

604 This is the fallback function if no version specific alteration are 

605 necessary. Effectively, a no-op. 

606 

607 Parameters 

608 ---------- 

609 **kwargs 

610 Any keyword arguments that Schedd.history(), Schedd.query(), and 

611 Schedd.xquery() accepts. 

612 

613 Returns 

614 ------- 

615 kwargs : `dict` [`str`, `~typing.Any`] 

616 Keywords arguments that were passed to the function. 

617 """ 

618 return kwargs 

619 

620 

621def htc_query_history(schedds, **kwargs): 

622 """Fetch history records from the condor_schedd daemon. 

623 

624 Parameters 

625 ---------- 

626 schedds : `htcondor.Schedd` 

627 HTCondor schedulers which to query for job information. 

628 **kwargs 

629 Any keyword arguments that Schedd.history() accepts. 

630 

631 Yields 

632 ------ 

633 schedd_name : `str` 

634 Name of the HTCondor scheduler managing the job queue. 

635 job_ad : `dict` [`str`, `~typing.Any`] 

636 A dictionary representing HTCondor ClassAd describing a job. It maps 

637 job attributes names to values of the ClassAd expressions they 

638 represent. 

639 """ 

640 # If not set, provide defaults for positional arguments. 

641 kwargs.setdefault("constraint", None) 

642 kwargs.setdefault("projection", []) 

643 kwargs = htc_tune_schedd_args(**kwargs) 

644 for schedd_name, schedd in schedds.items(): 

645 for job_ad in schedd.history(**kwargs): 

646 yield schedd_name, dict(job_ad) 

647 

648 

649def htc_query_present(schedds, **kwargs): 

650 """Query the condor_schedd daemon for job ads. 

651 

652 Parameters 

653 ---------- 

654 schedds : `htcondor.Schedd` 

655 HTCondor schedulers which to query for job information. 

656 **kwargs 

657 Any keyword arguments that Schedd.xquery() accepts. 

658 

659 Yields 

660 ------ 

661 schedd_name : `str` 

662 Name of the HTCondor scheduler managing the job queue. 

663 job_ad : `dict` [`str`, `~typing.Any`] 

664 A dictionary representing HTCondor ClassAd describing a job. It maps 

665 job attributes names to values of the ClassAd expressions they 

666 represent. 

667 """ 

668 kwargs = htc_tune_schedd_args(**kwargs) 

669 for schedd_name, schedd in schedds.items(): 

670 for job_ad in schedd.query(**kwargs): 

671 yield schedd_name, dict(job_ad) 

672 

673 

674def htc_version(): 

675 """Return the version given by the HTCondor API. 

676 

677 Returns 

678 ------- 

679 version : `str` 

680 HTCondor version as easily comparable string. 

681 """ 

682 return str(HTC_VERSION) 

683 

684 

685def htc_submit_dag(sub): 

686 """Submit job for execution. 

687 

688 Parameters 

689 ---------- 

690 sub : `htcondor.Submit` 

691 An object representing a job submit description. 

692 

693 Returns 

694 ------- 

695 schedd_job_info : `dict` [`str`, `dict` [`str`, \ 

696 `dict` [`str`, `~typing.Any`]]] 

697 Information about jobs satisfying the search criteria where for each 

698 Scheduler, local HTCondor job ids are mapped to their respective 

699 classads. 

700 """ 

701 coll = htcondor.Collector() 

702 schedd_ad = coll.locate(htcondor.DaemonTypes.Schedd) 

703 schedd = htcondor.Schedd(schedd_ad) 

704 

705 # If Schedd.submit() fails, the method will raise an exception. Usually, 

706 # that implies issues with the HTCondor pool which BPS can't address. 

707 # Hence, no effort is made to handle the exception. 

708 submit_result = schedd.submit(sub) 

709 

710 # Sadly, the ClassAd from Schedd.submit() (see above) does not have 

711 # 'GlobalJobId' so we need to run a regular query to get it anyway. 

712 schedd_name = schedd_ad["Name"] 

713 schedd_dag_info = condor_q( 

714 constraint=f"ClusterId == {submit_result.cluster()}", schedds={schedd_name: schedd} 

715 ) 

716 return schedd_dag_info 

717 

718 

719def htc_create_submit_from_dag( 

720 dag_filename: str, submit_options: dict[str, Any], dagman_conf_filename: str | os.PathLike | None = None 

721) -> htcondor.Submit: 

722 """Create a DAGMan job submit description. 

723 

724 Parameters 

725 ---------- 

726 dag_filename : `str` 

727 Name of file containing HTCondor DAG commands. 

728 submit_options : `dict` [`str`, `~typing.Any`], optional 

729 Contains extra options for command line (Value of None means flag). 

730 dagman_conf_filename : `str` or `os.PathLike`, optional 

731 Location of DAGMan configuration file, if Any. Defaults to no 

732 file (None). 

733 

734 Returns 

735 ------- 

736 sub : `htcondor.Submit` 

737 An object representing a job submit description. 

738 

739 Notes 

740 ----- 

741 Use with HTCondor versions which support htcondor.Submit.from_dag(), 

742 i.e., 8.9.3 or newer. 

743 """ 

744 # Config and environment variables do not seem to override -MaxIdle 

745 # on the .dag.condor.sub's command line (broken in some 24.0.x versions). 

746 # Explicitly forward them as a submit_option if either exists. 

747 # Note: auto generated subdag submit files are still the -MaxIdle=1000 

748 # in the broken versions. 

749 if "MaxIdle" not in submit_options: 

750 max_jobs_idle: int | None = None 

751 config_var_name = "DAGMAN_MAX_JOBS_IDLE" 

752 

753 if dagman_conf_filename: 

754 _LOG.debug("Checking DAGMan config file = %s", dagman_conf_filename) 

755 with open(dagman_conf_filename) as fh: 

756 for line in fh: 

757 _LOG.debug("DAGMan config file line = %s", line) 

758 parts = line.split("=") 

759 _LOG.debug("DAGMan config file line parts = %s", parts) 

760 if len(parts) == 2 and parts[0].strip() == config_var_name: 

761 max_jobs_idle = int(parts[1].strip()) 

762 _LOG.debug("Found %s = %s", config_var_name, max_jobs_idle) 

763 break 

764 if max_jobs_idle is None: 

765 if f"_CONDOR_{config_var_name}" in os.environ: 

766 max_jobs_idle = int(os.environ[f"_CONDOR_{config_var_name}"]) 

767 elif config_var_name in htcondor.param: 

768 max_jobs_idle = htcondor.param[config_var_name] 

769 

770 if max_jobs_idle: 

771 submit_options["MaxIdle"] = max_jobs_idle 

772 else: 

773 _LOG.debug("MaxIdle already in submit_options: %s", submit_options) 

774 

775 _LOG.debug("Using submit_options = %s", submit_options) 

776 return htcondor.Submit.from_dag(dag_filename, submit_options) 

777 

778 

779def htc_create_submit_from_cmd(dag_filename, submit_options=None): 

780 """Create a DAGMan job submit description. 

781 

782 Create a DAGMan job submit description by calling ``condor_submit_dag`` 

783 on given DAG description file. 

784 

785 Parameters 

786 ---------- 

787 dag_filename : `str` 

788 Name of file containing HTCondor DAG commands. 

789 submit_options : `dict` [`str`, `~typing.Any`], optional 

790 Contains extra options for command line (Value of None means flag). 

791 

792 Returns 

793 ------- 

794 sub : `htcondor.Submit` 

795 An object representing a job submit description. 

796 

797 Notes 

798 ----- 

799 Use with HTCondor versions which do not support htcondor.Submit.from_dag(), 

800 i.e., older than 8.9.3. 

801 """ 

802 # Run command line condor_submit_dag command. 

803 cmd = "condor_submit_dag -f -no_submit -notification never -autorescue 1 -UseDagDir -no_recurse " 

804 

805 if submit_options is not None: 

806 for opt, val in submit_options.items(): 

807 cmd += f" -{opt} {val or ''}" 

808 cmd += f"{dag_filename}" 

809 

810 process = subprocess.Popen( 

811 cmd.split(), shell=False, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, encoding="utf-8" 

812 ) 

813 process.wait() 

814 

815 if process.returncode != 0: 

816 print(f"Exit code: {process.returncode}") 

817 print(process.communicate()[0]) 

818 raise RuntimeError("Problems running condor_submit_dag") 

819 

820 return htc_create_submit_from_file(f"{dag_filename}.condor.sub") 

821 

822 

823def htc_create_submit_from_file(submit_file): 

824 """Parse a submission file. 

825 

826 Parameters 

827 ---------- 

828 submit_file : `str` 

829 Name of the HTCondor submit file. 

830 

831 Returns 

832 ------- 

833 sub : `htcondor.Submit` 

834 An object representing a job submit description. 

835 """ 

836 descriptors = {} 

837 with open(submit_file) as fh: 

838 for line in fh: 

839 line = line.strip() 

840 if not line.startswith("#") and not line == "queue": 

841 (key, val) = re.split(r"\s*=\s*", line, maxsplit=1) 

842 descriptors[key] = val 

843 

844 # Avoid UserWarning: the line 'copy_to_spool = False' was 

845 # unused by Submit object. Is it a typo? 

846 try: 

847 del descriptors["copy_to_spool"] 

848 except KeyError: 

849 pass 

850 

851 return htcondor.Submit(descriptors) 

852 

853 

854def _htc_write_job_commands(stream, name, commands, node_type="JOB"): 

855 """Output the DAGMan job lines for single job in DAG. 

856 

857 Parameters 

858 ---------- 

859 stream : `~typing.TextIO` 

860 Writeable text stream (typically an opened file). 

861 name : `str` 

862 Job name. 

863 commands : `RestrictedDict` 

864 DAG commands for a job. 

865 node_type : `str`, optional 

866 Type of DAGMan node (JOB, FINAL, SERVICE). Defaults to "JOB". 

867 """ 

868 # Note: optional pieces of commands include a space at the beginning. 

869 # also making sure values aren't empty strings as placeholders. 

870 if "pre" in commands and commands["pre"]: 

871 defer = "" 

872 if "defer" in commands["pre"] and commands["pre"]["defer"]: 

873 defer = f" DEFER {commands['pre']['defer']['status']} {commands['pre']['defer']['time']}" 

874 

875 debug = "" 

876 if "debug" in commands["pre"] and commands["pre"]["debug"]: 

877 debug = f" DEBUG {commands['pre']['debug']['filename']} {commands['pre']['debug']['type']}" 

878 

879 arguments = "" 

880 if "arguments" in commands["pre"] and commands["pre"]["arguments"]: 

881 arguments = f" {commands['pre']['arguments']}" 

882 

883 executable = commands["pre"]["executable"] 

884 print(f"SCRIPT{defer}{debug} PRE {name} {executable}{arguments}", file=stream) 

885 

886 if "post" in commands and commands["post"]: 

887 defer = "" 

888 if "defer" in commands["post"] and commands["post"]["defer"]: 

889 defer = f" DEFER {commands['post']['defer']['status']} {commands['post']['defer']['time']}" 

890 

891 debug = "" 

892 if "debug" in commands["post"] and commands["post"]["debug"]: 

893 debug = f" DEBUG {commands['post']['debug']['filename']} {commands['post']['debug']['type']}" 

894 

895 arguments = "" 

896 if "arguments" in commands["post"] and commands["post"]["arguments"]: 

897 arguments = f" {commands['post']['arguments']}" 

898 

899 executable = commands["post"]["executable"] 

900 print(f"SCRIPT{defer}{debug} POST {name} {executable}{arguments}", file=stream) 

901 

902 if "vars" in commands and commands["vars"]: 

903 for key, value in commands["vars"].items(): 

904 print(f'VARS {name} {key}="{htc_escape(value)}"', file=stream) 

905 

906 if "pre_skip" in commands and commands["pre_skip"]: 

907 print(f"PRE_SKIP {name} {commands['pre_skip']}", file=stream) 

908 

909 # FINAL node cannot have a DAGMan retry, abort-dag-on, priority, category 

910 if node_type != "FINAL": 

911 if "retry" in commands and commands["retry"]: 

912 print(f"RETRY {name} {commands['retry']}", end="", file=stream) 

913 if "retry_unless_exit" in commands and commands["retry_unless_exit"]: 

914 print(f" UNLESS-EXIT {commands['retry_unless_exit']}", end="", file=stream) 

915 print("", file=stream) # Since previous prints don't include new line 

916 

917 if "abort_dag_on" in commands and commands["abort_dag_on"]: 

918 print( 

919 f"ABORT-DAG-ON {name} {commands['abort_dag_on']['node_exit']}" 

920 f" RETURN {commands['abort_dag_on']['abort_exit']}", 

921 file=stream, 

922 ) 

923 

924 if "priority" in commands and commands["priority"]: 

925 print( 

926 f"PRIORITY {name} {commands['priority']}", 

927 file=stream, 

928 ) 

929 

930 

931class HTCJob: 

932 """HTCondor job for use in building DAG. 

933 

934 Parameters 

935 ---------- 

936 name : `str` 

937 Name of the job. 

938 label : `str` 

939 Label that can used for grouping or lookup. 

940 initcmds : `RestrictedDict` 

941 Initial job commands for submit file. 

942 initdagcmds : `RestrictedDict` 

943 Initial commands for job inside DAG. 

944 initattrs : `dict` 

945 Initial dictionary of job attributes. 

946 """ 

947 

948 def __init__(self, name, label=None, initcmds=(), initdagcmds=(), initattrs=None): 

949 self.name = name 

950 self.label = label 

951 self.cmds = RestrictedDict(HTC_VALID_JOB_KEYS, initcmds) 

952 self.dagcmds = RestrictedDict(HTC_VALID_JOB_DAG_KEYS, initdagcmds) 

953 self.attrs = initattrs 

954 self.subfile = None 

955 self.subdir = None 

956 self.subdag = None 

957 

958 def __str__(self): 

959 return self.name 

960 

961 def add_job_cmds(self, new_commands): 

962 """Add commands to Job (overwrite existing). 

963 

964 Parameters 

965 ---------- 

966 new_commands : `dict` 

967 Submit file commands to be added to Job. 

968 """ 

969 self.cmds.update(new_commands) 

970 

971 def add_dag_cmds(self, new_commands): 

972 """Add DAG commands to Job (overwrite existing). 

973 

974 Parameters 

975 ---------- 

976 new_commands : `dict` 

977 DAG file commands to be added to Job. 

978 """ 

979 self.dagcmds.update(new_commands) 

980 

981 def add_job_attrs(self, new_attrs): 

982 """Add attributes to Job (overwrite existing). 

983 

984 Parameters 

985 ---------- 

986 new_attrs : `dict` 

987 Attributes to be added to Job. 

988 """ 

989 if self.attrs is None: 

990 self.attrs = {} 

991 if new_attrs: 

992 self.attrs.update(new_attrs) 

993 

994 def write_submit_file(self, submit_path: str | os.PathLike) -> None: 

995 """Write job description to submit file. 

996 

997 Parameters 

998 ---------- 

999 submit_path : `str` or `os.PathLike` 

1000 Prefix path for the submit file. 

1001 """ 

1002 if not self.subfile: 

1003 self.subfile = f"{self.name}.sub" 

1004 

1005 subfile = self.subfile 

1006 if self.subdir: 

1007 subfile = Path(self.subdir) / subfile 

1008 

1009 subfile = Path(os.path.expandvars(subfile)) 

1010 if not subfile.is_absolute(): 

1011 subfile = Path(submit_path) / subfile 

1012 if not subfile.exists(): 

1013 _LOG.debug("Writing subfile: %s", subfile) 

1014 htc_write_condor_file(subfile, self.name, self.cmds, self.attrs) 

1015 else: 

1016 _LOG.debug("Using existing subfile: %s", subfile) 

1017 

1018 def write_dag_commands(self, stream, dag_rel_path, command_name="JOB"): 

1019 """Write DAG commands for single job to output stream. 

1020 

1021 Parameters 

1022 ---------- 

1023 stream : `~typing.TextIO` 

1024 Output Stream. 

1025 dag_rel_path : `str` 

1026 Relative path of dag to submit directory. 

1027 command_name : `str` 

1028 Name of the DAG command (e.g., JOB, FINAL). 

1029 """ 

1030 subfile = os.path.expandvars(self.subfile) 

1031 

1032 # JOB NodeName SubmitDescription [DIR directory] [NOOP] [DONE] 

1033 job_line = f'{command_name} {self.name} "{subfile}"' 

1034 if "dir" in self.dagcmds: 

1035 dir_val = self.dagcmds["dir"] 

1036 if dag_rel_path: 

1037 dir_val = os.path.join(dag_rel_path, dir_val) 

1038 job_line += f' DIR "{dir_val}"' 

1039 if self.dagcmds.get("noop", False): 

1040 job_line += " NOOP" 

1041 

1042 print(job_line, file=stream) 

1043 if self.dagcmds: 

1044 _htc_write_job_commands(stream, self.name, self.dagcmds, command_name) 

1045 

1046 def dump(self, fh): 

1047 """Dump job information to output stream. 

1048 

1049 Parameters 

1050 ---------- 

1051 fh : `~typing.TextIO` 

1052 Output stream. 

1053 """ 

1054 printer = pprint.PrettyPrinter(indent=4, stream=fh) 

1055 printer.pprint(self.name) 

1056 printer.pprint(self.cmds) 

1057 printer.pprint(self.attrs) 

1058 

1059 

1060class HTCDag(networkx.DiGraph): 

1061 """HTCondor DAG. 

1062 

1063 Parameters 

1064 ---------- 

1065 data : `~typing.Any` 

1066 Initial graph data of any format that is supported 

1067 by the to_network_graph() function. 

1068 name : `str` 

1069 Name for DAG. 

1070 """ 

1071 

1072 def __init__(self, data=None, name=""): 

1073 super().__init__(data=data, name=name) 

1074 

1075 self.graph["attr"] = {} 

1076 self.graph["run_id"] = None 

1077 self.graph["submit_path"] = None 

1078 self.graph["final_job"] = None 

1079 self.graph["service_job"] = None 

1080 self.graph["submit_options"] = {} 

1081 

1082 def __str__(self): 

1083 """Represent basic DAG info as string. 

1084 

1085 Returns 

1086 ------- 

1087 info : `str` 

1088 String containing basic DAG info. 

1089 """ 

1090 return f"{self.graph['name']} {len(self)}" 

1091 

1092 def add_attribs(self, attribs=None): 

1093 """Add attributes to the DAG. 

1094 

1095 Parameters 

1096 ---------- 

1097 attribs : `dict` 

1098 DAG attributes. 

1099 """ 

1100 if attribs is not None: 1100 ↛ exitline 1100 didn't return from function 'add_attribs' because the condition on line 1100 was always true

1101 self.graph["attr"].update(attribs) 

1102 

1103 def add_job(self, job, parent_names=None, child_names=None): 

1104 """Add an HTCJob to the HTCDag. 

1105 

1106 Parameters 

1107 ---------- 

1108 job : `HTCJob` 

1109 HTCJob to add to the HTCDag. 

1110 parent_names : `~collections.abc.Iterable` [`str`], optional 

1111 Names of parent jobs. 

1112 child_names : `~collections.abc.Iterable` [`str`], optional 

1113 Names of child jobs. 

1114 """ 

1115 assert isinstance(job, HTCJob) 

1116 _LOG.debug("Adding job %s to dag", job.name) 

1117 

1118 # Add dag level attributes to each job 

1119 job.add_job_attrs(self.graph["attr"]) 

1120 

1121 self.add_node(job.name, data=job) 

1122 

1123 if parent_names is not None: 1123 ↛ 1124line 1123 didn't jump to line 1124 because the condition on line 1123 was never true

1124 self.add_job_relationships(parent_names, [job.name]) 

1125 

1126 if child_names is not None: 1126 ↛ 1127line 1126 didn't jump to line 1127 because the condition on line 1126 was never true

1127 self.add_job_relationships(child_names, [job.name]) 

1128 

1129 def add_job_relationships(self, parents, children): 

1130 """Add DAG edge between parents and children jobs. 

1131 

1132 Parameters 

1133 ---------- 

1134 parents : `list` [`str`] 

1135 Contains parent job name(s). 

1136 children : `list` [`str`] 

1137 Contains children job name(s). 

1138 """ 

1139 self.add_edges_from(itertools.product(parents, children)) 

1140 

1141 def add_final_job(self, job): 

1142 """Add an HTCJob for the FINAL job in HTCDag. 

1143 

1144 Parameters 

1145 ---------- 

1146 job : `HTCJob` 

1147 HTCJob to add to the HTCDag as a FINAL job. 

1148 """ 

1149 # Add dag level attributes to each job 

1150 job.add_job_attrs(self.graph["attr"]) 

1151 

1152 self.graph["final_job"] = job 

1153 

1154 def add_service_job(self, job): 

1155 """Add an HTCJob for the SERVICE job in HTCDag. 

1156 

1157 Parameters 

1158 ---------- 

1159 job : `HTCJob` 

1160 HTCJob to add to the HTCDag as a SERVICE job. 

1161 """ 

1162 # Add dag level attributes to each job 

1163 job.add_job_attrs(self.graph["attr"]) 

1164 

1165 self.graph["service_job"] = job 

1166 

1167 def del_job(self, job_name): 

1168 """Delete the job from the DAG. 

1169 

1170 Parameters 

1171 ---------- 

1172 job_name : `str` 

1173 Name of job in DAG to delete. 

1174 """ 

1175 # Reconnect edges around node to delete 

1176 parents = self.predecessors(job_name) 

1177 children = self.successors(job_name) 

1178 self.add_edges_from(itertools.product(parents, children)) 

1179 

1180 # Delete job node (which deletes its edges). 

1181 self.remove_node(job_name) 

1182 

1183 def write(self, submit_path, job_subdir="", dag_subdir="", dag_rel_path=""): 

1184 """Write DAG to a file. 

1185 

1186 Parameters 

1187 ---------- 

1188 submit_path : `str` 

1189 Prefix path for all outputs. 

1190 job_subdir : `str`, optional 

1191 Template for job subdir (submit_path + job_subdir). 

1192 dag_subdir : `str`, optional 

1193 DAG subdir (submit_path + dag_subdir). 

1194 dag_rel_path : `str`, optional 

1195 Prefix to job_subdir for jobs inside subdag. 

1196 """ 

1197 self.graph["submit_path"] = submit_path 

1198 self.graph["dag_filename"] = os.path.join(dag_subdir, f"{self.graph['name']}.dag") 

1199 full_filename = os.path.join(submit_path, self.graph["dag_filename"]) 

1200 os.makedirs(os.path.dirname(full_filename), exist_ok=True) 

1201 

1202 try: 

1203 dagman_config_path = Path(self.graph["attr"]["bps_wms_config_path"]) 

1204 except KeyError: 

1205 dagman_config_path = None 

1206 with open(full_filename, "w") as fh: 

1207 if dagman_config_path is not None: 

1208 fh.write(f"CONFIG {dag_rel_path / dagman_config_path}\n") 

1209 

1210 for name, nodeval in self.nodes().items(): 

1211 try: 

1212 job = nodeval["data"] 

1213 except KeyError: 

1214 _LOG.error("Job %s doesn't have data (keys: %s).", name, nodeval.keys()) 

1215 raise 

1216 if job.subdag: 1216 ↛ 1217line 1216 didn't jump to line 1217 because the condition on line 1216 was never true

1217 if job.subfile: 

1218 this_dag_rel_path = "" 

1219 else: 

1220 this_dag_rel_path = "../.." 

1221 dag_subdir = f"subdags/{job.name}" 

1222 if "dir" in job.dagcmds: 

1223 subdir = job.dagcmds["dir"] 

1224 else: 

1225 subdir = job_subdir 

1226 if dagman_config_path is not None: 

1227 job.subdag.add_attribs({"bps_wms_config_path": str(dagman_config_path)}) 

1228 job.subdag.write(submit_path, subdir, dag_subdir, this_dag_rel_path) 

1229 fh.write(f"SUBDAG EXTERNAL {job.name} {Path(job.subdag.graph['dag_filename']).name}") 

1230 if dag_subdir: 

1231 fh.write(f" DIR {dag_subdir}") 

1232 fh.write("\n") 

1233 if job.dagcmds: 

1234 _htc_write_job_commands(fh, job.name, job.dagcmds) 

1235 else: 

1236 job.write_submit_file(submit_path) 

1237 job.write_dag_commands(fh, dag_rel_path) 

1238 

1239 for edge in self.edges(): 

1240 print(f"PARENT {edge[0]} CHILD {edge[1]}", file=fh) 

1241 

1242 if self.graph.get("write_dot", False): 

1243 print(f"DOT {self.name}.dot", file=fh) 

1244 

1245 print(f"NODE_STATUS_FILE {self.name}.node_status", file=fh) 

1246 

1247 # Add bps attributes to dag submission 

1248 for key, value in self.graph["attr"].items(): 

1249 print(f'SET_JOB_ATTR {key}= "{htc_escape(value)}"', file=fh) 

1250 

1251 # Add special nodes if any. 

1252 special_jobs = { 

1253 "FINAL": self.graph["final_job"], 

1254 "SERVICE": self.graph["service_job"], 

1255 } 

1256 for dagcmd, job in special_jobs.items(): 

1257 if job is not None: 

1258 job.write_submit_file(submit_path) 

1259 job.write_dag_commands(fh, dag_rel_path, dagcmd) 

1260 

1261 def dump(self, fh): 

1262 """Dump DAG info to output stream. 

1263 

1264 Parameters 

1265 ---------- 

1266 fh : `typing.IO` 

1267 Where to dump DAG info as text. 

1268 """ 

1269 for key, value in self.graph: 

1270 print(f"{key}={value}", file=fh) 

1271 for name, data in self.nodes().items(): 

1272 print(f"{name}:", file=fh) 

1273 data.dump(fh) 

1274 for edge in self.edges(): 

1275 print(f"PARENT {edge[0]} CHILD {edge[1]}", file=fh) 

1276 if self.graph["final_job"]: 

1277 print(f"FINAL {self.graph['final_job'].name}:", file=fh) 

1278 self.graph["final_job"].dump(fh) 

1279 

1280 def write_dot(self, filename): 

1281 """Write a dot version of the DAG. 

1282 

1283 Parameters 

1284 ---------- 

1285 filename : `str` 

1286 Name of the dot file. 

1287 """ 

1288 pos = networkx.nx_agraph.graphviz_layout(self) 

1289 networkx.draw(self, pos=pos) 

1290 networkx.drawing.nx_pydot.write_dot(self, filename) 

1291 

1292 

1293def condor_q(constraint=None, schedds=None, **kwargs): 

1294 """Get information about the jobs in the HTCondor job queue(s). 

1295 

1296 Parameters 

1297 ---------- 

1298 constraint : `str`, optional 

1299 Constraints to be passed to job query. 

1300 schedds : `dict` [`str`, `htcondor.Schedd`], optional 

1301 HTCondor schedulers which to query for job information. If None 

1302 (default), the query will be run against local scheduler only. 

1303 **kwargs : `~typing.Any` 

1304 Additional keyword arguments that need to be passed to the internal 

1305 query method. 

1306 

1307 Returns 

1308 ------- 

1309 job_info : `dict` [`str`, `dict` [`str`, `dict` [`str`, `~typing.Any`]]] 

1310 Information about jobs satisfying the search criteria where for each 

1311 Scheduler, local HTCondor job ids are mapped to their respective 

1312 classads. 

1313 """ 

1314 return condor_query(constraint, schedds, htc_query_present, **kwargs) 

1315 

1316 

1317def condor_history(constraint=None, schedds=None, **kwargs): 

1318 """Get information about the jobs from HTCondor history records. 

1319 

1320 Parameters 

1321 ---------- 

1322 constraint : `str`, optional 

1323 Constraints to be passed to job query. 

1324 schedds : `dict` [`str`, `htcondor.Schedd`], optional 

1325 HTCondor schedulers which to query for job information. If None 

1326 (default), the query will be run against the history file of 

1327 the local scheduler only. 

1328 **kwargs : `~typing.Any` 

1329 Additional keyword arguments that need to be passed to the internal 

1330 query method. 

1331 

1332 Returns 

1333 ------- 

1334 job_info : `dict` [`str`, `dict` [`str`, `dict` [`str`, `~typing.Any`]]] 

1335 Information about jobs satisfying the search criteria where for each 

1336 Scheduler, local HTCondor job ids are mapped to their respective 

1337 classads. 

1338 """ 

1339 return condor_query(constraint, schedds, htc_query_history, **kwargs) 

1340 

1341 

1342def condor_query(constraint=None, schedds=None, query_func=htc_query_present, **kwargs): 

1343 """Get information about HTCondor jobs. 

1344 

1345 Parameters 

1346 ---------- 

1347 constraint : `str`, optional 

1348 Constraints to be passed to job query. 

1349 schedds : `dict` [`str`, `htcondor.Schedd`], optional 

1350 HTCondor schedulers which to query for job information. If None 

1351 (default), the query will be run against the history file of 

1352 the local scheduler only. 

1353 query_func : `~collections.abc.Callable` 

1354 An query function which takes following arguments: 

1355 

1356 - ``schedds``: Schedulers to query (`list` [`htcondor.Schedd`]). 

1357 - ``**kwargs``: Keyword arguments that will be passed to the query 

1358 function. 

1359 **kwargs : `~typing.Any` 

1360 Additional keyword arguments that need to be passed to the query 

1361 method. 

1362 

1363 Returns 

1364 ------- 

1365 job_info : `dict` [`str`, `dict` [`str`, `dict` [`str`, `~typing.Any`]]] 

1366 Information about jobs satisfying the search criteria where for each 

1367 Scheduler, local HTCondor job ids are mapped to their respective 

1368 classads. 

1369 """ 

1370 if not schedds: 

1371 coll = htcondor.Collector() 

1372 schedd_ad = coll.locate(htcondor.DaemonTypes.Schedd) 

1373 schedds = {schedd_ad["Name"]: htcondor.Schedd(schedd_ad)} 

1374 

1375 # Make sure that 'ClusterId' and 'ProcId' attributes are always included 

1376 # in the job classad. They are needed to construct the job id. 

1377 added_attrs = set() 

1378 if "projection" in kwargs and kwargs["projection"]: 

1379 requested_attrs = set(kwargs["projection"]) 

1380 required_attrs = {"ClusterId", "ProcId"} 

1381 added_attrs = required_attrs - requested_attrs 

1382 for attr in added_attrs: 

1383 kwargs["projection"].append(attr) 

1384 

1385 unwanted_attrs = {"Env", "Environment"} | added_attrs 

1386 job_info = defaultdict(dict) 

1387 for schedd_name, job_ad in query_func(schedds, constraint=constraint, **kwargs): 

1388 id_ = f"{job_ad['ClusterId']}.{job_ad['ProcId']}" 

1389 for attr in set(job_ad) & unwanted_attrs: 

1390 del job_ad[attr] 

1391 job_info[schedd_name][id_] = job_ad 

1392 _LOG.debug("query returned %d jobs", sum(len(val) for val in job_info.values())) 

1393 

1394 # Restore the list of the requested attributes to its original value 

1395 # if needed. 

1396 if added_attrs: 

1397 for attr in added_attrs: 

1398 kwargs["projection"].remove(attr) 

1399 

1400 # When returning the results filter out entries for schedulers with no jobs 

1401 # matching the search criteria. 

1402 return {key: val for key, val in job_info.items() if val} 

1403 

1404 

1405def condor_search(constraint=None, hist=None, schedds=None): 

1406 """Search for running and finished jobs satisfying given criteria. 

1407 

1408 Parameters 

1409 ---------- 

1410 constraint : `str`, optional 

1411 Constraints to be passed to job query. 

1412 hist : `float` 

1413 Limit history search to this many days. 

1414 schedds : `dict` [`str`, `htcondor.Schedd`], optional 

1415 The list of the HTCondor schedulers which to query for job information. 

1416 If None (default), only the local scheduler will be queried. 

1417 

1418 Returns 

1419 ------- 

1420 job_info : `dict` [`str`, `dict` [`str`, `dict` [`str` `~typing.Any`]]] 

1421 Information about jobs satisfying the search criteria where for each 

1422 Scheduler, local HTCondor job ids are mapped to their respective 

1423 classads. 

1424 """ 

1425 if not schedds: 

1426 coll = htcondor.Collector() 

1427 schedd_ad = coll.locate(htcondor.DaemonTypes.Schedd) 

1428 schedds = {schedd_ad["Name"]: htcondor.Schedd(locate_ad=schedd_ad)} 

1429 

1430 job_info = condor_q(constraint=constraint, schedds=schedds) 

1431 if hist is not None: 

1432 _LOG.debug("Searching history going back %s days", hist) 

1433 epoch = (datetime.now() - timedelta(days=hist)).timestamp() 

1434 constraint += f" && (CompletionDate >= {epoch} || JobFinishedHookDone >= {epoch})" 

1435 hist_info = condor_history(constraint, schedds=schedds) 

1436 update_job_info(job_info, hist_info) 

1437 return job_info 

1438 

1439 

1440def condor_status(constraint=None, coll=None): 

1441 """Get information about HTCondor pool. 

1442 

1443 Parameters 

1444 ---------- 

1445 constraint : `str`, optional 

1446 Constraints to be passed to the query. 

1447 coll : `htcondor.Collector`, optional 

1448 Object representing HTCondor collector daemon. 

1449 

1450 Returns 

1451 ------- 

1452 pool_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

1453 Mapping between HTCondor slot names and slot information (classAds). 

1454 """ 

1455 if coll is None: 

1456 coll = htcondor.Collector() 

1457 try: 

1458 pool_ads = coll.query(constraint=constraint) 

1459 except OSError as ex: 

1460 raise RuntimeError(f"Problem querying the Collector. (Constraint='{constraint}')") from ex 

1461 

1462 pool_info = {} 

1463 for slot in pool_ads: 

1464 pool_info[slot["name"]] = dict(slot) 

1465 _LOG.debug("condor_status returned %d ads", len(pool_info)) 

1466 return pool_info 

1467 

1468 

1469def update_job_info(job_info, other_info): 

1470 """Update results of a job query with results from another query. 

1471 

1472 Parameters 

1473 ---------- 

1474 job_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

1475 Results of the job query that needs to be updated. 

1476 other_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

1477 Results of the other job query. 

1478 

1479 Returns 

1480 ------- 

1481 job_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

1482 The updated results. 

1483 """ 

1484 for schedd_name, others in other_info.items(): 

1485 try: 

1486 jobs = job_info[schedd_name] 

1487 except KeyError: 

1488 job_info[schedd_name] = others 

1489 else: 

1490 for id_, ad in others.items(): 

1491 jobs.setdefault(id_, {}).update(ad) 

1492 return job_info 

1493 

1494 

1495def count_jobs_in_single_dag( 

1496 filename: str | os.PathLike, 

1497) -> tuple[Counter[str], dict[str, str], dict[str, WmsNodeType]]: 

1498 """Build bps_run_summary string from dag file. 

1499 

1500 Parameters 

1501 ---------- 

1502 filename : `str` 

1503 Path that includes dag file for a run. 

1504 

1505 Returns 

1506 ------- 

1507 counts : `Counter` [`str`] 

1508 Semi-colon separated list of job labels and counts. 

1509 (Same format as saved in dag classad). 

1510 job_name_to_label : `dict` [`str`, `str`] 

1511 Mapping of job names to job labels. 

1512 job_name_to_type : `dict` [`str`, `lsst.ctrl.bps.htcondor.WmsNodeType`] 

1513 Mapping of job names to job types 

1514 (e.g., payload, final, service). 

1515 """ 

1516 # Later code depends upon insertion order 

1517 counts: Counter = Counter() # counts of payload jobs per label 

1518 job_name_to_label: dict[str, str] = {} 

1519 job_name_to_type: dict[str, WmsNodeType] = {} 

1520 with open(filename) as fh: 

1521 for line in fh: 

1522 # Skip any line that contains commands irrelevant to job counting. 

1523 if not line.startswith( 

1524 ( 

1525 "JOB", 

1526 "FINAL", 

1527 "SERVICE", 

1528 "SUBDAG EXTERNAL", 

1529 ) 

1530 ): 

1531 continue 

1532 

1533 m = re.match( 

1534 r"(?P<command>JOB|FINAL|SERVICE|SUBDAG EXTERNAL)\s+" 

1535 r'(?P<jobname>(?P<wms>wms_)?\S+)\s+"?(?P<subfile>\S+)"?\s*' 

1536 r'(DIR "?(?P<dir>[^\s"]+)"?)?\s*(?P<noop>NOOP)?', 

1537 line, 

1538 ) 

1539 if m: 1539 ↛ 1592line 1539 didn't jump to line 1592 because the condition on line 1539 was always true

1540 job_name = m.group("jobname") 

1541 job_type = WmsNodeType.UNKNOWN 

1542 name_parts = job_name.split("_") 

1543 

1544 label = "" 

1545 if m.group("dir"): 

1546 dir_match = re.search(r"jobs/([^\s/]+)", m.group("dir")) 

1547 if dir_match: 

1548 label = dir_match.group(1) 

1549 else: 

1550 _LOG.debug("Parse DAG: unparsed dir = %s", line) 

1551 elif m.group("subfile"): 1551 ↛ 1558line 1551 didn't jump to line 1558 because the condition on line 1551 was always true

1552 subfile_match = re.search(r"jobs/([^\s/]+)", m.group("subfile")) 

1553 if subfile_match: 

1554 label = m.group("subfile").split("/")[1] 

1555 else: 

1556 label = pegasus_name_to_label(job_name) 

1557 

1558 match m.group("command"): 

1559 case "JOB": 

1560 if m.group("noop"): 

1561 job_type = WmsNodeType.NOOP 

1562 # wms_noop_label 

1563 label = name_parts[2] 

1564 elif m.group("wms"): 

1565 if name_parts[1] == "check": 1565 ↛ 1570line 1565 didn't jump to line 1570 because the condition on line 1565 was always true

1566 job_type = WmsNodeType.SUBDAG_CHECK 

1567 # wms_check_status_wms_group_label 

1568 label = name_parts[5] 

1569 else: 

1570 _LOG.warning( 

1571 "Unexpected skipping of dag line due to unknown wms job: %s", line 

1572 ) 

1573 else: 

1574 job_type = WmsNodeType.PAYLOAD 

1575 if label == "init": 1575 ↛ 1576line 1575 didn't jump to line 1576 because the condition on line 1575 was never true

1576 label = "pipetaskInit" 

1577 counts[label] += 1 

1578 case "FINAL": 

1579 job_type = WmsNodeType.FINAL 

1580 counts[label] += 1 # final counts a payload job. 

1581 case "SERVICE": 

1582 job_type = WmsNodeType.SERVICE 

1583 case "SUBDAG EXTERNAL": 1583 ↛ 1587line 1583 didn't jump to line 1587 because the pattern on line 1583 always matched

1584 job_type = WmsNodeType.SUBDAG 

1585 label = name_parts[2] 

1586 

1587 job_name_to_label[job_name] = label 

1588 job_name_to_type[job_name] = job_type 

1589 else: 

1590 # The line should, but didn't match the pattern above. Probably 

1591 # problems with regex. 

1592 _LOG.warning("Unexpected skipping of dag line: %s", line) 

1593 

1594 return counts, job_name_to_label, job_name_to_type 

1595 

1596 

1597def summarize_dag(dir_name: str) -> tuple[str, dict[str, str], dict[str, WmsNodeType]]: 

1598 """Build bps_run_summary string from dag file. 

1599 

1600 Parameters 

1601 ---------- 

1602 dir_name : `str` 

1603 Path that includes dag file for a run. 

1604 

1605 Returns 

1606 ------- 

1607 summary : `str` 

1608 Semi-colon separated list of job labels and counts 

1609 (Same format as saved in dag classad). 

1610 job_name_to_label : `dict` [`str`, `str`] 

1611 Mapping of job names to job labels. 

1612 job_name_to_type : `dict` [`str`, `lsst.ctrl.bps.htcondor.WmsNodeType`] 

1613 Mapping of job names to job types 

1614 (e.g., payload, final, service). 

1615 """ 

1616 # Later code depends upon insertion order 

1617 counts: Counter[str] = Counter() # counts of payload jobs per label 

1618 job_name_to_label: dict[str, str] = {} 

1619 job_name_to_type: dict[str, WmsNodeType] = {} 

1620 for filename in Path(dir_name).glob("*.dag"): 

1621 single_counts, single_job_name_to_label, single_job_name_to_type = count_jobs_in_single_dag(filename) 

1622 counts += single_counts 

1623 _update_dicts(job_name_to_label, single_job_name_to_label) 

1624 _update_dicts(job_name_to_type, single_job_name_to_type) 

1625 

1626 for filename in Path(dir_name).glob("subdags/*/*.dag"): 

1627 single_counts, single_job_name_to_label, single_job_name_to_type = count_jobs_in_single_dag(filename) 

1628 counts += single_counts 

1629 _update_dicts(job_name_to_label, single_job_name_to_label) 

1630 _update_dicts(job_name_to_type, single_job_name_to_type) 

1631 

1632 summary = ";".join([f"{name}:{counts[name]}" for name in counts]) 

1633 _LOG.debug("summarize_dag: %s %s %s", summary, job_name_to_label, job_name_to_type) 

1634 return summary, job_name_to_label, job_name_to_type 

1635 

1636 

1637def pegasus_name_to_label(name): 

1638 """Convert pegasus job name to a label for the report. 

1639 

1640 Parameters 

1641 ---------- 

1642 name : `str` 

1643 Name of job. 

1644 

1645 Returns 

1646 ------- 

1647 label : `str` 

1648 Label for job. 

1649 """ 

1650 label = "UNK" 

1651 if name.startswith("create_dir") or name.startswith("stage_in") or name.startswith("stage_out"): 1651 ↛ 1652line 1651 didn't jump to line 1652 because the condition on line 1651 was never true

1652 label = "pegasus" 

1653 else: 

1654 m = re.match(r"pipetask_(\d+_)?([^_]+)", name) 

1655 if m: 1655 ↛ 1656line 1655 didn't jump to line 1656 because the condition on line 1655 was never true

1656 label = m.group(2) 

1657 if label == "init": 

1658 label = "pipetaskInit" 

1659 

1660 return label 

1661 

1662 

1663def read_single_dag_status(filename: str | os.PathLike) -> dict[str, Any]: 

1664 """Read the node status file for DAG summary information. 

1665 

1666 Parameters 

1667 ---------- 

1668 filename : `str` or `Path.pathlib` 

1669 Node status filename. 

1670 

1671 Returns 

1672 ------- 

1673 dag_ad : `dict` [`str`, `~typing.Any`] 

1674 DAG summary information. 

1675 """ 

1676 dag_ad: dict[str, Any] = {} 

1677 

1678 # While this is probably more up to date than dag classad, only read from 

1679 # file if need to. 

1680 try: 

1681 node_stat_file = Path(filename) 

1682 _LOG.debug("Reading Node Status File %s", node_stat_file) 

1683 with open(node_stat_file) as infh: 

1684 dag_ad = dict(classad.parseNext(infh)) # pylint: disable=E1101 

1685 

1686 if not dag_ad: 1686 ↛ 1688line 1686 didn't jump to line 1688 because the condition on line 1686 was never true

1687 # Pegasus check here 

1688 metrics_file = node_stat_file.with_suffix(".dag.metrics") 

1689 if metrics_file.exists(): 

1690 with open(metrics_file) as infh: 

1691 metrics = json.load(infh) 

1692 dag_ad["NodesTotal"] = metrics.get("jobs", 0) 

1693 dag_ad["NodesFailed"] = metrics.get("jobs_failed", 0) 

1694 dag_ad["NodesDone"] = metrics.get("jobs_succeeded", 0) 

1695 metrics_file = node_stat_file.with_suffix(".metrics") 

1696 with open(metrics_file) as infh: 

1697 metrics = json.load(infh) 

1698 dag_ad["NodesTotal"] = metrics["wf_metrics"]["total_jobs"] 

1699 except (OSError, PermissionError): 

1700 pass 

1701 

1702 _LOG.debug("read_dag_status: %s", dag_ad) 

1703 return dag_ad 

1704 

1705 

1706def read_dag_status(wms_path: str | os.PathLike) -> dict[str, Any]: 

1707 """Read the node status file for DAG summary information. 

1708 

1709 Parameters 

1710 ---------- 

1711 wms_path : `str` or `os.PathLike` 

1712 Path that includes node status file for a run. 

1713 

1714 Returns 

1715 ------- 

1716 dag_ad : `dict` [`str`, `~typing.Any`] 

1717 DAG summary information, counts summed across any subdags. 

1718 """ 

1719 dag_ads: dict[str, Any] = {} 

1720 path = Path(wms_path) 

1721 try: 

1722 node_stat_file = next(path.glob("*.node_status")) 

1723 except StopIteration as exc: 

1724 raise FileNotFoundError(f"DAGMan node status not found in {wms_path}") from exc 

1725 

1726 dag_ads = read_single_dag_status(node_stat_file) 

1727 

1728 for node_stat_file in path.glob("subdags/*/*.node_status"): 

1729 dag_ad = read_single_dag_status(node_stat_file) 

1730 dag_ads["JobProcsHeld"] += dag_ad.get("JobProcsHeld", 0) 

1731 dag_ads["NodesPost"] += dag_ad.get("NodesPost", 0) 

1732 dag_ads["JobProcsIdle"] += dag_ad.get("JobProcsIdle", 0) 

1733 dag_ads["NodesTotal"] += dag_ad.get("NodesTotal", 0) 

1734 dag_ads["NodesFailed"] += dag_ad.get("NodesFailed", 0) 

1735 dag_ads["NodesDone"] += dag_ad.get("NodesDone", 0) 

1736 dag_ads["NodesQueued"] += dag_ad.get("NodesQueued", 0) 

1737 dag_ads["NodesPre"] += dag_ad.get("NodesReady", 0) 

1738 dag_ads["NodesFutile"] += dag_ad.get("NodesFutile", 0) 

1739 dag_ads["NodesUnready"] += dag_ad.get("NodesUnready", 0) 

1740 

1741 return dag_ads 

1742 

1743 

1744def read_single_node_status(filename: str | os.PathLike, init_fake_id: int) -> dict[str, Any]: 

1745 """Read entire node status file. 

1746 

1747 Parameters 

1748 ---------- 

1749 filename : `str` or `pathlib.Path` 

1750 Node status filename. 

1751 init_fake_id : `int` 

1752 Initial fake id value. 

1753 

1754 Returns 

1755 ------- 

1756 jobs : `dict` [`str`, `~typing.Any`] 

1757 DAG summary information compiled from the node status file combined 

1758 with the information found in the node event log. 

1759 

1760 Currently, if the same job attribute is found in both files, its value 

1761 from the event log takes precedence over the value from the node status 

1762 file. 

1763 """ 

1764 filename = Path(filename) 

1765 

1766 # Get jobid info from other places to fill in gaps in info from node_status 

1767 _, job_name_to_label, job_name_to_type = count_jobs_in_single_dag(filename.with_suffix(".dag")) 

1768 loginfo: dict[str, dict[str, Any]] = {} 

1769 wms_workflow_id = MISSING_ID 

1770 try: 

1771 wms_workflow_id, _ = read_single_dag_log(filename.with_suffix(".dag.dagman.log")) 

1772 loginfo = read_single_dag_nodes_log(filename.with_suffix(".dag.nodes.log")) 

1773 except (OSError, PermissionError): 

1774 pass 

1775 

1776 job_name_to_id: dict[str, str] = {} 

1777 _LOG.debug("loginfo = %s", loginfo) 

1778 log_job_name_to_id: dict[str, str] = {} 

1779 for job_id, job_info in loginfo.items(): 

1780 if "LogNotes" in job_info: 1780 ↛ 1779line 1780 didn't jump to line 1779 because the condition on line 1780 was always true

1781 m = re.match(r"DAG Node: (\S+)", job_info["LogNotes"]) 

1782 if m: 1782 ↛ 1779line 1782 didn't jump to line 1779 because the condition on line 1782 was always true

1783 job_name = m.group(1) 

1784 log_job_name_to_id[job_name] = job_id 

1785 job_info["DAGNodeName"] = job_name 

1786 job_info["wms_node_type"] = job_name_to_type[job_name] 

1787 job_info["bps_job_label"] = job_name_to_label[job_name] 

1788 

1789 jobs = {} 

1790 fake_id = init_fake_id # For nodes that do not yet have a job id, give fake one 

1791 try: 

1792 with open(filename) as fh: 

1793 for ad in classad.parseAds(fh): 

1794 match ad["Type"]: 

1795 case "DagStatus": 

1796 # Skip DAG summary. 

1797 pass 

1798 case "NodeStatus": 

1799 job_name = ad["Node"] 

1800 if job_name in job_name_to_label: 1800 ↛ 1802line 1800 didn't jump to line 1802 because the condition on line 1800 was always true

1801 job_label = job_name_to_label[job_name] 

1802 elif "_" in job_name: 

1803 job_label = job_name.split("_")[1] 

1804 else: 

1805 job_label = job_name 

1806 

1807 job = dict(ad) 

1808 if job_name in log_job_name_to_id: 

1809 job_id = str(log_job_name_to_id[job_name]) 

1810 _update_dicts(job, loginfo[job_id]) 

1811 else: 

1812 job_id = str(fake_id) 

1813 job = dict(ad) 

1814 fake_id -= 1 

1815 jobs[job_id] = job 

1816 job_name_to_id[job_name] = job_id 

1817 

1818 # Make job info as if came from condor_q. 

1819 job["ClusterId"] = int(float(job_id)) 

1820 job["DAGManJobID"] = wms_workflow_id 

1821 job["DAGNodeName"] = job_name 

1822 job["bps_job_label"] = job_label 

1823 job["wms_node_type"] = job_name_to_type[job_name] 

1824 

1825 case "StatusEnd": 1825 ↛ 1828line 1825 didn't jump to line 1828 because the pattern on line 1825 always matched

1826 # Skip node status file "epilog". 

1827 pass 

1828 case _: 

1829 _LOG.debug( 

1830 "Ignoring unknown classad type '%s' in the node status file '%s'", 

1831 ad["Type"], 

1832 filename, 

1833 ) 

1834 except (OSError, PermissionError): 

1835 pass 

1836 

1837 # Check for missing jobs (e.g., submission failure or not submitted yet) 

1838 # Use dag info to create job placeholders 

1839 for name in set(job_name_to_label) - set(job_name_to_id): 

1840 if name in log_job_name_to_id: # job was in nodes.log, but not node_status 

1841 job_id = str(log_job_name_to_id[name]) 

1842 job = dict(loginfo[job_id]) 

1843 else: 

1844 job_id = str(fake_id) 

1845 fake_id -= 1 

1846 job = {} 

1847 job["NodeStatus"] = NodeStatus.NOT_READY 

1848 

1849 job["ClusterId"] = int(float(job_id)) 

1850 job["ProcId"] = 0 

1851 job["DAGManJobID"] = wms_workflow_id 

1852 job["DAGNodeName"] = name 

1853 job["bps_job_label"] = job_name_to_label[name] 

1854 job["wms_node_type"] = job_name_to_type[name] 

1855 jobs[f"{job['ClusterId']}.{job['ProcId']}"] = job 

1856 

1857 for job_info in jobs.values(): 

1858 job_info["from_dag_job"] = f"wms_{filename.stem}" 

1859 

1860 return jobs 

1861 

1862 

1863def read_node_status(wms_path: str | os.PathLike) -> dict[str, dict[str, Any]]: 

1864 """Read entire node status file. 

1865 

1866 Parameters 

1867 ---------- 

1868 wms_path : `str` or `os.PathLike` 

1869 Path that includes node status file for a run. 

1870 

1871 Returns 

1872 ------- 

1873 jobs : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

1874 DAG summary information compiled from the node status file combined 

1875 with the information found in the node event log. 

1876 

1877 Currently, if the same job attribute is found in both files, its value 

1878 from the event log takes precedence over the value from the node status 

1879 file. 

1880 """ 

1881 jobs: dict[str, dict[str, Any]] = {} 

1882 init_fake_id = -1 

1883 

1884 # subdags may not have run so wouldn't have node_status file 

1885 # use dag files and let read_single_node_status handle missing 

1886 # node_status file. 

1887 for dag_filename in Path(wms_path).glob("*.dag"): 

1888 filename = dag_filename.with_suffix(".node_status") 

1889 info = read_single_node_status(filename, init_fake_id) 

1890 init_fake_id -= len(info) 

1891 _update_dicts(jobs, info) 

1892 

1893 for dag_filename in Path(wms_path).glob("subdags/*/*.dag"): 

1894 filename = dag_filename.with_suffix(".node_status") 

1895 info = read_single_node_status(filename, init_fake_id) 

1896 init_fake_id -= len(info) 

1897 _update_dicts(jobs, info) 

1898 

1899 # Propagate pruned from subdags to jobs 

1900 name_to_id: dict[str, str] = {} 

1901 missing_status: dict[str, list[str]] = {} 

1902 for id_, job in jobs.items(): 

1903 if job["DAGNodeName"].startswith("wms_"): 

1904 name_to_id[job["DAGNodeName"]] = id_ 

1905 if "NodeStatus" not in job or job["NodeStatus"] == NodeStatus.NOT_READY: 

1906 missing_status.setdefault(job["from_dag_job"], []).append(id_) 

1907 

1908 for name, dag_id in name_to_id.items(): 

1909 dag_status = jobs[dag_id].get("NodeStatus", NodeStatus.NOT_READY) 

1910 if dag_status in {NodeStatus.NOT_READY, NodeStatus.FUTILE}: 

1911 for id_ in missing_status.get(name, []): 

1912 jobs[id_]["NodeStatus"] = dag_status 

1913 

1914 return jobs 

1915 

1916 

1917def read_single_dag_log(log_filename: str | os.PathLike) -> tuple[str, dict[str, dict[str, Any]]]: 

1918 """Read job information from the DAGMan log file. 

1919 

1920 Parameters 

1921 ---------- 

1922 log_filename : `str` or `os.PathLike` 

1923 DAGMan log filename. 

1924 

1925 Returns 

1926 ------- 

1927 wms_workflow_id : `str` 

1928 HTCondor job id (i.e., <ClusterId>.<ProcId>) of the DAGMan job. 

1929 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

1930 HTCondor job information read from the log file mapped to HTCondor 

1931 job id. 

1932 

1933 Raises 

1934 ------ 

1935 FileNotFoundError 

1936 If cannot find DAGMan log in given wms_path. 

1937 """ 

1938 wms_workflow_id = MISSING_ID 

1939 dag_info: dict[str, dict[str, Any]] = {} 

1940 

1941 filename = Path(log_filename) 

1942 if filename.exists(): 

1943 _LOG.debug("dag node log filename: %s", filename) 

1944 

1945 info: dict[str, Any] = {} 

1946 job_event_log = htcondor.JobEventLog(str(filename)) 

1947 for event in job_event_log.events(stop_after=0): 

1948 id_ = f"{event['Cluster']}.{event['Proc']}" 

1949 if id_ not in info: 

1950 info[id_] = {} 

1951 wms_workflow_id = id_ # taking last job id in case of restarts 

1952 _update_dicts(info[id_], event) 

1953 info[id_][f"{event.type.name.lower()}_time"] = event["EventTime"] 

1954 

1955 # only save latest DAG job 

1956 dag_info = {wms_workflow_id: info[wms_workflow_id]} 

1957 

1958 return wms_workflow_id, dag_info 

1959 

1960 

1961def read_dag_log(wms_path: str | os.PathLike) -> tuple[str, dict[str, Any]]: 

1962 """Read job information from the DAGMan log file. 

1963 

1964 Parameters 

1965 ---------- 

1966 wms_path : `str` or `os.PathLike` 

1967 Path containing the DAGMan log file. 

1968 

1969 Returns 

1970 ------- 

1971 wms_workflow_id : `str` 

1972 HTCondor job id (i.e., <ClusterId>.<ProcId>) of the DAGMan job. 

1973 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

1974 HTCondor job information read from the log file mapped to HTCondor 

1975 job id. 

1976 

1977 Raises 

1978 ------ 

1979 FileNotFoundError 

1980 If cannot find DAGMan log in given wms_path. 

1981 """ 

1982 wms_workflow_id = MISSING_ID 

1983 dag_info: dict[str, dict[str, Any]] = {} 

1984 

1985 path = Path(wms_path) 

1986 if path.exists(): 1986 ↛ 1998line 1986 didn't jump to line 1998 because the condition on line 1986 was always true

1987 # Can be more than one dag file in directory. Assume one 

1988 # with lowest ID is main DAG 

1989 ids = [] 

1990 for filename in path.glob("*.dag.dagman.log"): 

1991 _LOG.debug("dag log filename: %s", filename) 

1992 single_id, single_dag_info = read_single_dag_log(filename) 

1993 _update_dicts(dag_info, single_dag_info) 

1994 ids.append(single_id) 

1995 if ids: 

1996 wms_workflow_id = min(ids) 

1997 

1998 if wms_workflow_id == MISSING_ID: 

1999 raise FileNotFoundError(f"DAGMan log not found in {wms_path}") 

2000 

2001 return wms_workflow_id, dag_info 

2002 

2003 

2004def read_single_dag_nodes_log(filename: str | os.PathLike) -> dict[str, dict[str, Any]]: 

2005 """Read job information from the DAGMan nodes log file. 

2006 

2007 Parameters 

2008 ---------- 

2009 filename : `str` or `os.PathLike` 

2010 Path containing the DAGMan nodes log file. 

2011 

2012 Returns 

2013 ------- 

2014 info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

2015 HTCondor job information read from the log file mapped to HTCondor 

2016 job id. 

2017 

2018 Raises 

2019 ------ 

2020 FileNotFoundError 

2021 If cannot find DAGMan node log in given wms_path. 

2022 """ 

2023 _LOG.debug("dag node log filename: %s", filename) 

2024 filename = Path(filename) 

2025 

2026 info: dict[str, dict[str, Any]] = {} 

2027 if not filename.exists(): 

2028 raise FileNotFoundError(f"{filename} does not exist") 

2029 

2030 try: 

2031 job_event_log = htcondor.JobEventLog(str(filename)) 

2032 except htcondor.HTCondorIOError as ex: 

2033 _LOG.error("Problem reading nodes log file (%s): %s", filename, ex) 

2034 import traceback 

2035 

2036 traceback.print_stack() 

2037 raise 

2038 for event in job_event_log.events(stop_after=0): 

2039 _LOG.debug("log event type = %s, keys = %s", event["EventTypeNumber"], event.keys()) 

2040 

2041 try: 

2042 id_ = f"{event['Cluster']}.{event['Proc']}" 

2043 except KeyError: 

2044 _LOG.warn( 

2045 "Log event missing ids (DAGNodeName=%s, EventTime=%s, EventTypeNumber=%s)", 

2046 event.get("DAGNodeName", "UNK"), 

2047 event.get("EventTime", "UNK"), 

2048 event.get("EventTypeNumber", "UNK"), 

2049 ) 

2050 else: 

2051 if id_ not in info: 

2052 info[id_] = {} 

2053 # Workaround: Please check to see if still problem in 

2054 # future HTCondor versions. Sometimes get a 

2055 # JobAbortedEvent for a subdag job after it already 

2056 # terminated normally. Seems to happen when using job 

2057 # plus subdags. 

2058 if event["EventTypeNumber"] == 9 and info[id_].get("EventTypeNumber", -1) == 5: 2058 ↛ 2059line 2058 didn't jump to line 2059 because the condition on line 2058 was never true

2059 _LOG.debug("Skipping spurious JobAbortedEvent: %s", dict(event)) 

2060 elif event["EventTypeNumber"] == 16 and event["DAGNodeName"] == "finalJob": 

2061 # FINAL job's post script exit code is special and indicates 

2062 # status of DAG instead of just the FINAL job. Save the 

2063 # information separately. 

2064 info[id_]["post"] = dict(event) 

2065 else: 

2066 _update_dicts(info[id_], event) 

2067 info[id_][f"{event.type.name.lower()}_time"] = event["EventTime"] 

2068 

2069 return info 

2070 

2071 

2072def read_dag_nodes_log(wms_path: str | os.PathLike) -> dict[str, dict[str, Any]]: 

2073 """Read job information from the DAGMan nodes log file. 

2074 

2075 Parameters 

2076 ---------- 

2077 wms_path : `str` or `os.PathLike` 

2078 Path containing the DAGMan nodes log file. 

2079 

2080 Returns 

2081 ------- 

2082 info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

2083 HTCondor job information read from the log file mapped to HTCondor 

2084 job id. 

2085 

2086 Raises 

2087 ------ 

2088 FileNotFoundError 

2089 If cannot find DAGMan node log in given wms_path. 

2090 """ 

2091 info: dict[str, dict[str, Any]] = {} 

2092 for filename in Path(wms_path).glob("*.dag.nodes.log"): 

2093 _LOG.debug("dag node log filename: %s", filename) 

2094 _update_dicts(info, read_single_dag_nodes_log(filename)) 

2095 

2096 # If submitted, the main nodes log file should exist 

2097 if not info: 

2098 raise FileNotFoundError(f"DAGMan node log not found in {wms_path}") 

2099 

2100 # Subdags will not have dag nodes log files if they haven't 

2101 # started running yet (so missing is not an error). 

2102 for filename in Path(wms_path).glob("subdags/*/*.dag.nodes.log"): 

2103 _LOG.debug("dag node log filename: %s", filename) 

2104 _update_dicts(info, read_single_dag_nodes_log(filename)) 

2105 

2106 return info 

2107 

2108 

2109def read_dag_info(wms_path: str | os.PathLike) -> tuple[Path, dict[str, dict[str, Any]]]: 

2110 """Read custom DAGMan job information from the file. 

2111 

2112 Parameters 

2113 ---------- 

2114 wms_path : `str` or `os.PathLike` 

2115 Path containing the file with the DAGMan job info. 

2116 

2117 Returns 

2118 ------- 

2119 filename : `pathlib.Path` 

2120 Name of file containing the dag information. 

2121 dag_info : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

2122 HTCondor job information. 

2123 

2124 Raises 

2125 ------ 

2126 FileNotFoundError 

2127 If cannot find DAGMan job info file in the given location. 

2128 """ 

2129 dag_info: dict[str, dict[str, Any]] = {} 

2130 try: 

2131 filename = next(Path(wms_path).glob("*.info.json")) 

2132 except StopIteration as exc: 

2133 raise FileNotFoundError(f"File with DAGMan job information not found in {wms_path}") from exc 

2134 _LOG.debug("DAGMan job information filename: %s", filename) 

2135 try: 

2136 with open(filename) as fh: 

2137 dag_info = json.load(fh) 

2138 except (OSError, PermissionError) as exc: 

2139 _LOG.debug("Retrieving DAGMan job information failed: %s", exc) 

2140 return filename, dag_info 

2141 

2142 

2143def write_dag_info(filename: str, schedd_dag_info: dict[str, dict[str, Any]]): 

2144 """Write custom job information about DAGMan job. 

2145 

2146 Parameters 

2147 ---------- 

2148 filename : `str` 

2149 Name of the file where the information will be stored. If not given, 

2150 creates a filename using bps_run value. 

2151 schedd_dag_info : `dict` [`str` `dict` [`str`, `~typing.Any`]] 

2152 Information about the DAGMan job. 

2153 """ 

2154 _LOG.debug("schedd_dag_info = %s", schedd_dag_info) 

2155 schedd_name, dag_info = next(iter(schedd_dag_info.items())) 

2156 dag_id, dag_ad = next(iter(dag_info.items())) 

2157 ad = {"ClusterId": dag_ad["ClusterId"], "GlobalJobId": dag_ad["GlobalJobId"]} 

2158 ad.update({key: val for key, val in dag_ad.items() if key.startswith("bps")}) 

2159 try: 

2160 with open(filename, "w") as fh: 

2161 info = {schedd_name: {dag_id: ad}} 

2162 json.dump(info, fh) 

2163 except (KeyError, OSError, PermissionError) as exc: 

2164 _LOG.debug("Persisting DAGMan job information failed: %s", exc) 

2165 

2166 return filename 

2167 

2168 

2169def htc_tweak_log_info(wms_path: str | Path, job: dict[str, Any]) -> None: 

2170 """Massage the given job info has same structure as if came from condor_q. 

2171 

2172 Parameters 

2173 ---------- 

2174 wms_path : `str` | `os.PathLike` 

2175 Path containing an HTCondor event log file. 

2176 job : `dict` [ `str`, `~typing.Any` ] 

2177 A mapping between HTCondor job id and job information read from 

2178 the log. 

2179 """ 

2180 _LOG.debug("htc_tweak_log_info: %s %s", wms_path, job) 

2181 

2182 # Use the presence of 'MyType' key as a proxy to determine if the job ad 

2183 # contains the info extracted from the event log. Exit early if it doesn't 

2184 # (e.g. it is a job ad for a pruned job). 

2185 if "MyType" not in job: 

2186 return 

2187 

2188 try: 

2189 job["ClusterId"] = job["Cluster"] 

2190 job["ProcId"] = job["Proc"] 

2191 except KeyError as e: 

2192 _LOG.error("Missing key %s in job: %s", str(e), job) 

2193 raise 

2194 job["Iwd"] = str(wms_path) 

2195 job["Owner"] = Path(wms_path).owner() 

2196 

2197 if "LogNotes" in job: 

2198 m = re.match(r"DAG Node: (\S+)", job["LogNotes"]) 

2199 if m: 2199 ↛ 2202line 2199 didn't jump to line 2202 because the condition on line 2199 was always true

2200 job["DAGNodeName"] = m.group(1) 

2201 

2202 match job["MyType"]: 

2203 case "ExecuteEvent": 

2204 job["JobStatus"] = htcondor.JobStatus.RUNNING 

2205 case "JobTerminatedEvent" | "PostScriptTerminatedEvent": 

2206 job["JobStatus"] = htcondor.JobStatus.COMPLETED 

2207 case "SubmitEvent": 

2208 job["JobStatus"] = htcondor.JobStatus.IDLE 

2209 case "JobAbortedEvent": 

2210 job["JobStatus"] = htcondor.JobStatus.REMOVED 

2211 case "JobHeldEvent": 

2212 job["JobStatus"] = htcondor.JobStatus.HELD 

2213 case "JobReleaseEvent": 

2214 # If the job managing the execution of the root DAG is held and 

2215 # released this will be the last event showing up in its 

2216 # job event log even if the job is still running. If this is 

2217 # the last event for a job corresponding to the workflow node 

2218 # (either a normal payload job or the job managing the execution 

2219 # of an inner DAG), its final status will be determined later 

2220 # using node status log (see _htc_status_to_wms_state()). 

2221 job["JobStatus"] = htcondor.JobStatus.RUNNING if "DAGNodeName" not in job else None 

2222 case _: 

2223 _LOG.debug("Unknown log event type: %s", job["MyType"]) 

2224 job["JobStatus"] = None 

2225 

2226 # Use available information to add either "ExitCode" or "ExitSignal" 

2227 # attribute that captures respectively job's exit status (if it finished 

2228 # on its own accord) or its exit signal (if it was terminated by 

2229 # a signal). Also, include a flag "ExitBySignal" to make distinguishing 

2230 # between these two cases easy later on. 

2231 if job["JobStatus"] in { 

2232 htcondor.JobStatus.COMPLETED, 

2233 htcondor.JobStatus.HELD, 

2234 htcondor.JobStatus.REMOVED, 

2235 }: 

2236 new_job = HTC_JOB_AD_HANDLERS.handle(job) 

2237 if new_job is not None: 

2238 job = new_job 

2239 else: 

2240 _LOG.error("Could not determine exit status for job '%s.%s'", job["ClusterId"], job["ProcId"]) 

2241 

2242 

2243def htc_check_dagman_output(wms_path: str | os.PathLike) -> str: 

2244 """Check the DAGMan output for error messages. 

2245 

2246 Parameters 

2247 ---------- 

2248 wms_path : `str` or `os.PathLike` 

2249 Directory containing the DAGman output file. 

2250 

2251 Returns 

2252 ------- 

2253 message : `str` 

2254 Message containing error messages from the DAGMan output. Empty 

2255 string if no messages. 

2256 

2257 Raises 

2258 ------ 

2259 FileNotFoundError 

2260 If cannot find DAGMan standard output file in given wms_path. 

2261 """ 

2262 try: 

2263 filename = next(Path(wms_path).glob("*.dag.dagman.out")) 

2264 except StopIteration as exc: 

2265 raise FileNotFoundError(f"DAGMan standard output file not found in {wms_path}") from exc 

2266 _LOG.debug("dag output filename: %s", filename) 

2267 

2268 p = re.compile(r"^(\d\d/\d\d/\d\d \d\d:\d\d:\d\d) (Job submit try \d+/\d+ failed|Warning:.*$|ERROR:.*$)") 

2269 

2270 message = "" 

2271 # Shouldn't normally see this warning. Initializing it just in case. 

2272 last_warning = ( 

2273 "Missing warning: Usually means the .dagman.out file has been truncated " 

2274 "or temporarily had a full quota." 

2275 ) 

2276 try: 

2277 with open(filename) as fh: 

2278 last_submit_failed = "" # Since submit retries multiple times only report last one 

2279 for line in fh: 

2280 m = p.match(line) 

2281 if m: 

2282 if m.group(2).startswith("Job submit try"): 

2283 last_submit_failed = m.group(1) 

2284 elif m.group(2).startswith("ERROR: submit attempt failed"): 

2285 pass # Should be handled by Job submit try 

2286 elif m.group(2).startswith("Warning"): 

2287 if ".dag.nodes.log is in /tmp" in m.group(2): 

2288 last_warning = "Cannot submit from /tmp." 

2289 else: 

2290 last_warning = m.group(2) 

2291 elif m.group(2) == "ERROR: Warning is fatal error because of DAGMAN_USE_STRICT setting": 

2292 message += "ERROR: " 

2293 message += last_warning 

2294 message += "\n" 

2295 elif m.group(2) in [ 2295 ↛ 2301line 2295 didn't jump to line 2301 because the condition on line 2295 was always true

2296 "ERROR: the following job(s) failed:", 

2297 "ERROR: the following Node(s) failed:", 

2298 ]: 

2299 pass 

2300 else: 

2301 message += m.group(2) 

2302 message += "\n" 

2303 

2304 if last_submit_failed: 

2305 message += f"Warn: Job submission issues (last: {last_submit_failed})" 

2306 except (OSError, PermissionError): 

2307 message = f"Warn: Could not read dagman output file from {wms_path}." 

2308 _LOG.debug("dag output file message: %s", message) 

2309 return message 

2310 

2311 

2312def _read_rescue_headers(infh: TextIO) -> list[str]: 

2313 """Read header lines from a rescue file. 

2314 

2315 Parameters 

2316 ---------- 

2317 infh : `TextIO` 

2318 The rescue file from which to read the header lines. 

2319 

2320 Returns 

2321 ------- 

2322 header_lines : `list` [`str`] 

2323 Header lines read from the rescue file. 

2324 """ 

2325 header_lines = [] 

2326 for line in infh: 

2327 line = line.strip() 

2328 if not line.startswith("#"): 

2329 break 

2330 header_lines.append(line) 

2331 return header_lines 

2332 

2333 

2334def _update_rescue_headers(header_lines: list[str]) -> list[str]: 

2335 """Update rescue header lines in place. 

2336 

2337 Replaces ``wms_check_status`` node name prefixes with the corresponding 

2338 group node names and adjusts the count of nodes premarked DONE. 

2339 

2340 Parameters 

2341 ---------- 

2342 header_lines : `list` [`str`] 

2343 Header lines to update in place. 

2344 

2345 Returns 

2346 ------- 

2347 failed_subdags : `list` [`str`] 

2348 Names of failed subdag jobs. 

2349 """ 

2350 failed_subdags = [] 

2351 

2352 # Relies on HTCondor writing all failed node names on a single line 

2353 # (see src/condor_dagman/dag.cpp in https://github.com/htcondor/htcondor). 

2354 for i, line in enumerate(header_lines): 

2355 if line.startswith("# Nodes that failed:"): 

2356 nodes = header_lines[i + 1][1:].strip().split(",") 

2357 new_nodes = [] 

2358 for node in nodes: 

2359 if node.startswith("wms_check_status"): 

2360 node = node[17:] 

2361 failed_subdags.append(node) 

2362 new_nodes.append(node) 

2363 header_lines[i + 1] = f"# {','.join(new_nodes)}" 

2364 break 

2365 

2366 for i, line in enumerate(header_lines): 

2367 m = re.match(r"^(# Nodes premarked DONE:)\s+(\d+)", line) 

2368 if m: 

2369 header_lines[i] = f"{m.group(1)} {int(m.group(2)) - len(failed_subdags)}" 

2370 break 

2371 

2372 return failed_subdags 

2373 

2374 

2375def _write_rescue_headers(header_lines: list[str], outfh: TextIO) -> None: 

2376 """Write the header lines to the new rescue file. 

2377 

2378 Parameters 

2379 ---------- 

2380 header_lines : `list` [`str`] 

2381 Header lines to write to the new rescue file. 

2382 outfh : `TextIO` 

2383 New rescue file. 

2384 """ 

2385 for line in header_lines: 

2386 print(line, file=outfh) 

2387 print("", file=outfh) 

2388 

2389 

2390def _copy_done_lines(failed_subdags: list[str], infh: TextIO, outfh: TextIO) -> None: 

2391 """Copy the DONE lines from the original rescue file skipping 

2392 the failed group jobs. 

2393 

2394 Parameters 

2395 ---------- 

2396 failed_subdags : `list` [`str`] 

2397 List of job names for the failed subdags 

2398 infh : `TextIO` 

2399 Original rescue file to copy from. 

2400 outfh : `TextIO` 

2401 New rescue file to copy to. 

2402 """ 

2403 for line in infh: 

2404 line = line.strip() 

2405 try: 

2406 _, node_name = line.split() 

2407 except ValueError: 

2408 _LOG.error(f"Unexpected line in rescue file = '{line}'") 

2409 raise 

2410 if node_name not in failed_subdags: 

2411 print(line, file=outfh) 

2412 

2413 

2414def _update_rescue_file(rescue_file: Path) -> list[str]: 

2415 """Update the subdag failures in the main rescue file. 

2416 

2417 Parameters 

2418 ---------- 

2419 rescue_file : `pathlib.Path` 

2420 The main rescue file that needs to be updated. 

2421 

2422 Returns 

2423 ------- 

2424 failed_subdags : `list` [`str`] 

2425 Names of failed subdag jobs found in the DAGMan rescue file. 

2426 """ 

2427 # To reduce memory requirements, not reading entire file into memory. 

2428 rescue_tmp = rescue_file.with_suffix(rescue_file.suffix + ".tmp") 

2429 with open(rescue_file) as infh: 

2430 header_lines = _read_rescue_headers(infh) 

2431 failed_subdags = _update_rescue_headers(header_lines) 

2432 with open(rescue_tmp, "w") as outfh: 

2433 _write_rescue_headers(header_lines, outfh) 

2434 _copy_done_lines(failed_subdags, infh, outfh) 

2435 rescue_file.unlink() 

2436 rescue_tmp.rename(rescue_file) 

2437 return failed_subdags 

2438 

2439 

2440def _update_dicts(dict1, dict2): 

2441 """Update dict1 with info in dict2. 

2442 

2443 (Basically an update for nested dictionaries.) 

2444 

2445 Parameters 

2446 ---------- 

2447 dict1 : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

2448 HTCondor job information to be updated. 

2449 dict2 : `dict` [`str`, `dict` [`str`, `~typing.Any`]] 

2450 Additional HTCondor job information. 

2451 """ 

2452 for key, value in dict2.items(): 

2453 if key in dict1 and isinstance(dict1[key], dict) and isinstance(value, dict): 2453 ↛ 2454line 2453 didn't jump to line 2454 because the condition on line 2453 was never true

2454 _update_dicts(dict1[key], value) 

2455 else: 

2456 dict1[key] = value 

2457 

2458 

2459def _locate_schedds(locate_all=False): 

2460 """Find out Scheduler daemons in an HTCondor pool. 

2461 

2462 Parameters 

2463 ---------- 

2464 locate_all : `bool`, optional 

2465 If True, all available schedulers in the HTCondor pool will be located. 

2466 False by default which means that the search will be limited to looking 

2467 for the Scheduler running on a local host. 

2468 

2469 Returns 

2470 ------- 

2471 schedds : `dict` [`str`, `htcondor.Schedd`] 

2472 A mapping between Scheduler names and Python objects allowing for 

2473 interacting with them. 

2474 """ 

2475 coll = htcondor.Collector() 

2476 

2477 schedd_ads = [] 

2478 if locate_all: 

2479 schedd_ads.extend(coll.locateAll(htcondor.DaemonTypes.Schedd)) 

2480 else: 

2481 schedd_ads.append(coll.locate(htcondor.DaemonTypes.Schedd)) 

2482 return {ad["Name"]: htcondor.Schedd(ad) for ad in schedd_ads}