Coverage for python/lsst/ctrl/bps/drivers.py: 61%

255 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-27 09:23 +0000

1# This file is part of ctrl_bps. 

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 <http://www.gnu.org/licenses/>. 

27 

28"""Driver functions for each subcommand. 

29 

30Driver functions ensure that ensure all setup work is done before running 

31the subcommand method. 

32""" 

33 

34__all__ = [ 

35 "acquire_qgraph_driver", 

36 "batch_acquire_driver", 

37 "batch_prepare_driver", 

38 "cancel_driver", 

39 "cluster_qgraph_driver", 

40 "ping_driver", 

41 "prepare_driver", 

42 "report_driver", 

43 "restart_driver", 

44 "status_driver", 

45 "submit_driver", 

46 "submitcmd_driver", 

47 "transform_driver", 

48] 

49 

50 

51import logging 

52import os 

53from pathlib import Path 

54from typing import Any 

55 

56from lsst.pipe.base.quantum_graph import PredictedQuantumGraph 

57from lsst.resources import ResourcePath 

58from lsst.utils.timer import time_this 

59from lsst.utils.usage import get_peak_mem_usage 

60 

61from . import ( 

62 BPS_DEFAULTS, 

63 BPS_SEARCH_ORDER, 

64 DEFAULT_MEM_FMT, 

65 DEFAULT_MEM_UNIT, 

66 BpsConfig, 

67 ClusteredQuantumGraph, 

68 GenericWorkflow, 

69) 

70from .batch_submit import batch_payload_prepare, batch_submit 

71from .bps_reports import compile_code_summary, compile_job_summary 

72from .bps_utils import _dump_env_info, _dump_pkg_info, _make_id_link 

73from .cancel import cancel 

74from .construct import construct 

75from .initialize import ( 

76 custom_job_validator, 

77 init_submission, 

78 out_collection_validator, 

79 output_run_validator, 

80 submit_path_validator, 

81 translate_command_line_values, 

82) 

83from .ping import ping 

84from .pre_transform import acquire_quantum_graph, cluster_quanta, read_quantum_graph 

85from .prepare import prepare 

86from .report import display_report, retrieve_report 

87from .restart import restart 

88from .status import status 

89from .submit import submit 

90from .transform import transform 

91 

92_LOG = logging.getLogger(__name__) 

93 

94 

95def _init_submission_driver(config_file: str, **kwargs) -> BpsConfig: 

96 """Initialize runtime environment. 

97 

98 Parameters 

99 ---------- 

100 config_file : `str` 

101 Name of the configuration file. 

102 **kwargs 

103 Additional modifiers to the configuration. 

104 

105 Returns 

106 ------- 

107 config : `lsst.ctrl.bps.BpsConfig` 

108 Batch Processing Service configuration. 

109 """ 

110 validators = [submit_path_validator, output_run_validator, out_collection_validator] 

111 _LOG.info("Initializing BPS configuration and creating submit directory") 

112 with time_this( 

113 log=_LOG, 

114 level=logging.INFO, 

115 prefix=None, 

116 msg="BPS configuration initialized and submit directory created", 

117 mem_usage=True, 

118 mem_unit=DEFAULT_MEM_UNIT, 

119 mem_fmt=DEFAULT_MEM_FMT, 

120 ): 

121 config = init_submission(config_file, validators=validators, **kwargs) 

122 _log_mem_usage() 

123 

124 submit_path = config[".bps_defined.submitPath"] 

125 print(f"Submit dir: {submit_path}") 

126 return config 

127 

128 

129def acquire_qgraph_driver(config_file: str, **kwargs) -> tuple[BpsConfig, PredictedQuantumGraph]: 

130 """Read a quantum graph from a file or create one from pipeline definition. 

131 

132 Parameters 

133 ---------- 

134 config_file : `str` 

135 Name of the configuration file. 

136 **kwargs : `~typing.Any` 

137 Additional modifiers to the configuration. 

138 

139 Returns 

140 ------- 

141 config : `lsst.ctrl.bps.BpsConfig` 

142 Updated configuration. 

143 qgraph : `lsst.pipe.base.quantum_graph.PredictedQuantumGraph` 

144 A graph representing quanta. 

145 """ 

146 config = _init_submission_driver(config_file, **kwargs) 

147 

148 _LOG.info("Starting acquire stage (generating and/or reading quantum graph)") 

149 submit_path = config[".bps_defined.submitPath"] 

150 with time_this( 

151 log=_LOG, 

152 level=logging.INFO, 

153 prefix=None, 

154 msg="Acquire stage completed", 

155 mem_usage=True, 

156 mem_unit=DEFAULT_MEM_UNIT, 

157 mem_fmt=DEFAULT_MEM_FMT, 

158 ): 

159 qgraph_file = acquire_quantum_graph(config, out_prefix=submit_path) 

160 qgraph = read_quantum_graph(qgraph_file) 

161 

162 _log_mem_usage() 

163 

164 config[".bps_defined.runQgraphFile"] = qgraph_file 

165 return config, qgraph 

166 

167 

168def cluster_qgraph_driver(config_file: str, **kwargs: Any) -> tuple[BpsConfig, ClusteredQuantumGraph]: 

169 """Group quanta into clusters. 

170 

171 Parameters 

172 ---------- 

173 config_file : `str` 

174 Name of the configuration file. 

175 **kwargs : `~typing.Any` 

176 Additional modifiers to the configuration. 

177 

178 Returns 

179 ------- 

180 config : `lsst.ctrl.bps.BpsConfig` 

181 Updated configuration. 

182 clustered_qgraph : `lsst.ctrl.bps.ClusteredQuantumGraph` 

183 A graph representing clustered quanta. 

184 """ 

185 config, qgraph = acquire_qgraph_driver(config_file, **kwargs) 

186 

187 _LOG.info("Starting cluster stage (grouping quanta into jobs)") 

188 with time_this( 

189 log=_LOG, 

190 level=logging.INFO, 

191 prefix=None, 

192 msg="Cluster stage completed", 

193 mem_usage=True, 

194 mem_unit=DEFAULT_MEM_UNIT, 

195 mem_fmt=DEFAULT_MEM_FMT, 

196 ): 

197 clustered_qgraph = cluster_quanta(config, qgraph, config["uniqProcName"]) 

198 _log_mem_usage() 

199 

200 _LOG.info("ClusteredQuantumGraph contains %d cluster(s)", len(clustered_qgraph)) 

201 

202 submit_path = config[".bps_defined.submitPath"] 

203 _, save_clustered_qgraph = config.search("saveClusteredQgraph", opt={"default": False}) 

204 if save_clustered_qgraph: 

205 clustered_qgraph.save(os.path.join(submit_path, "bps_clustered_qgraph.pickle")) 

206 _, save_dot = config.search("saveDot", opt={"default": False}) 

207 if save_dot: 

208 clustered_qgraph.draw(os.path.join(submit_path, "bps_clustered_qgraph.dot")) 

209 return config, clustered_qgraph 

210 

211 

212def transform_driver(config_file: str, **kwargs: Any) -> tuple[BpsConfig, GenericWorkflow]: 

213 """Create a workflow for a specific workflow management system. 

214 

215 Parameters 

216 ---------- 

217 config_file : `str` 

218 Name of the configuration file. 

219 **kwargs : `~typing.Any` 

220 Additional modifiers to the configuration. 

221 

222 Returns 

223 ------- 

224 generic_workflow_config : `lsst.ctrl.bps.BpsConfig` 

225 Configuration to use when creating the workflow. 

226 generic_workflow : `lsst.ctrl.bps.GenericWorkflow` 

227 Representation of the abstract/scientific workflow specific to a given 

228 workflow management system. 

229 """ 

230 config, clustered_qgraph = cluster_qgraph_driver(config_file, **kwargs) 

231 submit_path = config[".bps_defined.submitPath"] 

232 

233 _LOG.info("Starting transform stage (creating generic workflow)") 

234 with time_this( 

235 log=_LOG, 

236 level=logging.INFO, 

237 prefix=None, 

238 msg="Transform stage completed", 

239 mem_usage=True, 

240 mem_unit=DEFAULT_MEM_UNIT, 

241 mem_fmt=DEFAULT_MEM_FMT, 

242 ): 

243 generic_workflow, generic_workflow_config = transform(config, clustered_qgraph, submit_path) 

244 _LOG.info("Generic workflow name '%s'", generic_workflow.name) 

245 _log_mem_usage() 

246 

247 num_jobs = sum(generic_workflow.job_counts.values()) 

248 _LOG.info("GenericWorkflow contains %d job(s) (including final)", num_jobs) 

249 

250 _, save_workflow = config.search("saveGenericWorkflow", opt={"default": False}) 

251 if save_workflow: 

252 with open(os.path.join(submit_path, "bps_generic_workflow.pickle"), "wb") as outfh: 

253 generic_workflow.save(outfh, "pickle") 

254 _, save_dot = config.search("saveDot", opt={"default": False}) 

255 if save_dot: 

256 with open(os.path.join(submit_path, "bps_generic_workflow.dot"), "w") as outfh: 

257 generic_workflow.draw(outfh, "dot") 

258 return generic_workflow_config, generic_workflow 

259 

260 

261def prepare_driver(config_file, **kwargs): 

262 """Create a representation of the generic workflow. 

263 

264 Parameters 

265 ---------- 

266 config_file : `str` 

267 Name of the configuration file. 

268 **kwargs : `~typing.Any` 

269 Additional modifiers to the configuration. 

270 

271 Returns 

272 ------- 

273 wms_config : `lsst.ctrl.bps.BpsConfig` 

274 Configuration to use when creating the workflow. 

275 workflow : `lsst.ctrl.bps.BaseWmsWorkflow` 

276 Representation of the abstract/scientific workflow specific to a given 

277 workflow management system. 

278 """ 

279 kwargs.setdefault("runWmsSubmissionChecks", True) 

280 generic_workflow_config, generic_workflow = transform_driver(config_file, **kwargs) 

281 submit_path = generic_workflow_config[".bps_defined.submitPath"] 

282 

283 _LOG.info("Starting prepare stage (creating specific implementation of workflow)") 

284 with time_this( 

285 log=_LOG, 

286 level=logging.INFO, 

287 prefix=None, 

288 msg="Prepare stage completed", 

289 mem_usage=True, 

290 mem_unit=DEFAULT_MEM_UNIT, 

291 mem_fmt=DEFAULT_MEM_FMT, 

292 ): 

293 wms_workflow = prepare(generic_workflow_config, generic_workflow, submit_path) 

294 _log_mem_usage() 

295 

296 wms_workflow_config = generic_workflow_config 

297 return wms_workflow_config, wms_workflow 

298 

299 

300def submit_driver(config_file, **kwargs): 

301 """Submit workflow for execution. 

302 

303 Parameters 

304 ---------- 

305 config_file : `str` 

306 Name of the configuration file. 

307 **kwargs : `~typing.Any` 

308 Additional modifiers to the configuration. 

309 """ 

310 kwargs.setdefault("runWmsSubmissionChecks", True) 

311 

312 _LOG.info( 

313 "DISCLAIMER: All values regarding memory consumption reported below are approximate and may " 

314 "not accurately reflect actual memory usage by the bps process." 

315 ) 

316 

317 config = BpsConfig( 

318 config_file, 

319 search_order=BPS_SEARCH_ORDER, 

320 defaults=BPS_DEFAULTS, 

321 wms_service_class_fqn=kwargs.get("wms_service"), 

322 ) 

323 translate_command_line_values(config, **kwargs) 

324 

325 wms_service_class = config["wmsServiceClass"] 

326 search_opts = config.get_search_opts() 

327 

328 # Initialization is normally called as part of submission stages. 

329 # But if running submission stages as batch job(s), need to 

330 # run initialization separately. 

331 

332 # PanDA-specific original options to run at sites with own Butler. 

333 search_opts["default"] = {} 

334 remote_build = config.search("remoteBuild", opt=search_opts)[1] 

335 remote_build_enabled = False 

336 if remote_build: # remoteBuild is a section of yaml 

337 search_opts["default"] = False 

338 remote_build_enabled = remote_build.search("enabled", opt=search_opts)[1] 

339 

340 # BPS option to turn on running submission stages as batch jobs 

341 search_opts["default"] = False 

342 batch_submission_enabled = config.search("bpsBatchSubmission", opt=search_opts)[1] 

343 

344 if remote_build_enabled or batch_submission_enabled: 

345 _LOG.info("Running submission stages as batch job(s) is enabled.") 

346 config = _init_submission_driver(config_file, **kwargs) 

347 

348 if wms_service_class == "lsst.ctrl.bps.panda.PanDAService": 348 ↛ 349line 348 didn't jump to line 349 because the condition on line 348 was never true

349 kwargs["remote_build"] = remote_build 

350 kwargs["config_file"] = config_file 

351 else: 

352 _LOG.info("The workflow is submitted to the local Data Facility.") 

353 

354 _LOG.info("Starting submission process") 

355 with time_this( 

356 log=_LOG, 

357 level=logging.INFO, 

358 prefix=None, 

359 msg="Submission process completed", 

360 mem_usage=True, 

361 mem_unit=DEFAULT_MEM_UNIT, 

362 mem_fmt=DEFAULT_MEM_FMT, 

363 ): 

364 if batch_submission_enabled: 

365 wms_workflow_config = config 

366 wms_workflow = batch_submit(config) 

367 else: 

368 if remote_build_enabled: 

369 wms_workflow_config = config 

370 wms_workflow = None 

371 else: 

372 wms_workflow_config, wms_workflow = prepare_driver(config_file, **kwargs) 

373 

374 _LOG.info("Starting submit stage") 

375 with time_this( 

376 log=_LOG, 

377 level=logging.INFO, 

378 prefix=None, 

379 msg="Submit stage completed", 

380 mem_usage=True, 

381 mem_unit=DEFAULT_MEM_UNIT, 

382 mem_fmt=DEFAULT_MEM_FMT, 

383 ): 

384 workflow = submit(wms_workflow_config, wms_workflow, **kwargs) 

385 if not wms_workflow: 

386 wms_workflow = workflow 

387 _LOG.info( 

388 "Run '%s' submitted for execution with id '%s'", wms_workflow.name, wms_workflow.run_id 

389 ) 

390 _log_mem_usage() 

391 

392 _make_id_link(wms_workflow_config, wms_workflow.run_id) 

393 

394 print(f"Run Id: {wms_workflow.run_id}") 

395 print(f"Run Name: {wms_workflow_config['uniqProcName']}") 

396 

397 

398def restart_driver(wms_service, run_id): 

399 """Restart a failed workflow. 

400 

401 Parameters 

402 ---------- 

403 wms_service : `str` 

404 Name of the class. 

405 run_id : `str` 

406 Id or path of workflow that need to be restarted. 

407 """ 

408 if wms_service is None: 

409 default_config = BpsConfig({}, defaults=BPS_DEFAULTS) 

410 wms_service = default_config["wmsServiceClass"] 

411 

412 new_run_id, run_name, message = restart(wms_service, run_id) 

413 if new_run_id is not None: 

414 path = Path(run_id) 

415 if path.exists(): 

416 _dump_env_info(f"{run_id}/{run_name}.env.info.yaml") 

417 _dump_pkg_info(f"{run_id}/{run_name}.pkg.info.yaml") 

418 config = BpsConfig(f"{run_id}/{run_name}_config.yaml") 

419 _make_id_link(config, new_run_id) 

420 

421 print(f"Run Id: {new_run_id}") 

422 print(f"Run Name: {run_name}") 

423 else: 

424 if message: 

425 print(f"Restart failed: {message}") 

426 else: 

427 print("Restart failed: Unknown error") 

428 

429 

430def report_driver( 

431 wms_service: str | None = None, 

432 run_id: str | None = None, 

433 user: str | None = None, 

434 hist_days: float = 0.0, 

435 pass_thru: str | None = None, 

436 is_global: bool = False, 

437 return_exit_codes: bool = False, 

438): 

439 """Print out the summary of jobs submitted for execution. 

440 

441 Parameters 

442 ---------- 

443 wms_service : `str`, optional 

444 Name of the class. 

445 run_id : `str`, optional 

446 A run id the report will be restricted to. 

447 user : `str`, optional 

448 A user the report will be restricted to. 

449 hist_days : `float`, optional 

450 Number of past days to consider while preparing the report. By default, 

451 only the currently running workflows are included in the report. 

452 If the report is restricted to a single run (i.e., ``run_id`` is set), 

453 the history search will be limited by default to two past days. 

454 pass_thru : `str`, optional 

455 A string to pass directly to the WMS service class. 

456 is_global : `bool`, optional 

457 If set, all available job queues will be queried for job information. 

458 Defaults to False which means that only a local job queue will be 

459 queried for information. 

460 

461 Only applicable in the context of a WMS using distributed job queues 

462 (e.g., HTCondor). 

463 return_exit_codes : `bool`, optional 

464 If set, return exit codes related to jobs with a 

465 non-success status. Defaults to False, which means that only 

466 the summary state is returned. 

467 

468 Only applicable in the context of a WMS with associated 

469 handlers to return exit codes from jobs. 

470 """ 

471 if not wms_service: 

472 default_config = BpsConfig(BPS_DEFAULTS) 

473 wms_service = os.environ.get("BPS_WMS_SERVICE_CLASS", default_config["wmsServiceClass"]) 

474 

475 # When reporting on a single run: 

476 # * increase history until a better mechanism for handling completed jobs 

477 # is available. 

478 # * massage the retrieved reports using BPS report postprocessors. 

479 if run_id: 

480 hist_days = max(hist_days, 2) 

481 postprocessors = [compile_job_summary] 

482 if return_exit_codes: 

483 postprocessors.append(compile_code_summary) 

484 else: 

485 postprocessors = None 

486 

487 runs, messages = retrieve_report( 

488 wms_service, 

489 run_id=run_id, 

490 user=user, 

491 hist=hist_days, 

492 pass_thru=pass_thru, 

493 is_global=is_global, 

494 return_exit_codes=return_exit_codes, 

495 postprocessors=postprocessors, 

496 ) 

497 

498 if runs or messages: 

499 display_report( 

500 runs, 

501 messages, 

502 is_detailed=bool(run_id), 

503 is_global=is_global, 

504 return_exit_codes=return_exit_codes, 

505 ) 

506 else: 

507 if run_id: 

508 print( 

509 f"No records found for job id '{run_id}'. " 

510 f"Hints: Double check id, retry with a larger --hist value (currently: {hist_days}), " 

511 "and/or use --global to search all job queues." 

512 ) 

513 

514 

515def status_driver(wms_service: str, run_id: str, hist_days: float, is_global: bool = False) -> int: 

516 """Print out status of workflow submitted for execution. 

517 

518 Parameters 

519 ---------- 

520 wms_service : `str` 

521 Name of the class. 

522 run_id : `str` 

523 A run id the report will be restricted to. 

524 hist_days : `float` 

525 Number of days. 

526 is_global : `bool`, optional 

527 If set, all available job queues will be queried for job information. 

528 Defaults to False which means that only a local job queue will be 

529 queried for information. 

530 

531 Only applicable in the context of a WMS using distributed job queues 

532 (e.g., HTCondor). 

533 

534 Returns 

535 ------- 

536 state : `int` 

537 Status of submitted workflow. 

538 """ 

539 if wms_service is None: 

540 default_config = BpsConfig(BPS_DEFAULTS) 

541 wms_service = os.environ.get("BPS_WMS_SERVICE_CLASS", default_config["wmsServiceClass"]) 

542 

543 state, message = status( 

544 wms_service, 

545 run_id=run_id, 

546 hist=hist_days, 

547 is_global=is_global, 

548 ) 

549 

550 _LOG.info("status: %s", state.name) 

551 if message: 

552 _LOG.warning(message) 

553 

554 return state.value 

555 

556 

557def cancel_driver(wms_service, run_id, user, require_bps, pass_thru, is_global=False): 

558 """Cancel submitted workflows. 

559 

560 Parameters 

561 ---------- 

562 wms_service : `str` 

563 Name of the Workload Management System service class. 

564 run_id : `str` 

565 ID or path of job that should be canceled. 

566 user : `str` 

567 User whose submitted jobs should be canceled. 

568 require_bps : `bool` 

569 Whether to require given run_id/user to be a bps submitted job. 

570 pass_thru : `str` 

571 Information to pass through to WMS. 

572 is_global : `bool`, optional 

573 If set, all available job queues will be checked for jobs to cancel. 

574 Defaults to False which means that only a local job queue will be 

575 checked. 

576 

577 Only applicable in the context of a WMS using distributed job queues 

578 (e.g., HTCondor). 

579 """ 

580 if wms_service is None: 

581 default_config = BpsConfig({}, defaults=BPS_DEFAULTS) 

582 wms_service = default_config["wmsServiceClass"] 

583 cancel(wms_service, run_id, user, require_bps, pass_thru, is_global=is_global) 

584 

585 

586def ping_driver(wms_service=None, pass_thru=None): 

587 """Check whether WMS services are up, reachable, and any authentication, 

588 if needed, succeeds. 

589 

590 The services to be checked are those needed for submit, report, cancel, 

591 restart, but ping cannot guarantee whether jobs would actually run 

592 successfully. 

593 

594 Parameters 

595 ---------- 

596 wms_service : `str`, optional 

597 Name of the Workload Management System service class. 

598 pass_thru : `str`, optional 

599 Information to pass through to WMS. 

600 

601 Returns 

602 ------- 

603 success : `int` 

604 Whether services are up and usable (0) or not (non-zero). 

605 """ 

606 if wms_service is None: 

607 default_config = BpsConfig({}, defaults=BPS_DEFAULTS) 

608 wms_service = default_config["wmsServiceClass"] 

609 status, message = ping(wms_service, pass_thru) 

610 

611 if message: 

612 if not status: 

613 _LOG.info(message) 

614 else: 

615 _LOG.error(message) 

616 

617 # Log overall status message 

618 if not status: 

619 _LOG.info("Ping successful.") 

620 else: 

621 _LOG.error("Ping failed (%d).", status) 

622 

623 return status 

624 

625 

626def submitcmd_driver(config_file: str, **kwargs) -> None: 

627 """Submit a command for execution. 

628 

629 Parameters 

630 ---------- 

631 config_file : `str` 

632 Name of the configuration file. 

633 **kwargs : `~typing.Any` 

634 Additional modifiers to the configuration. 

635 """ 

636 validators = [submit_path_validator, custom_job_validator] 

637 _LOG.info("Initializing BPS configuration and creating submit directory") 

638 with time_this( 

639 log=_LOG, 

640 level=logging.INFO, 

641 prefix=None, 

642 msg="BPS configuration initialized and submit directory created", 

643 mem_usage=True, 

644 mem_unit=DEFAULT_MEM_UNIT, 

645 mem_fmt=DEFAULT_MEM_FMT, 

646 ): 

647 config = init_submission(config_file, validators=validators, **kwargs) 

648 _log_mem_usage() 

649 

650 submit_path = config[".bps_defined.submitPath"] 

651 

652 _LOG.info("Starting construction stage (creating generic workflow)") 

653 with time_this( 

654 log=_LOG, 

655 level=logging.INFO, 

656 prefix=None, 

657 msg="Construction stage completed", 

658 mem_usage=True, 

659 mem_unit=DEFAULT_MEM_UNIT, 

660 mem_fmt=DEFAULT_MEM_FMT, 

661 ): 

662 generic_workflow, generic_workflow_config = construct(config) 

663 _LOG.info("Generic workflow name '%s'", generic_workflow.name) 

664 _log_mem_usage() 

665 

666 _, save_workflow = config.search("saveGenericWorkflow", opt={"default": False}) 

667 if save_workflow: 

668 with open(os.path.join(submit_path, "bps_generic_workflow.pickle"), "wb") as outfh: 

669 generic_workflow.save(outfh, "pickle") 

670 _, save_dot = config.search("saveDot", opt={"default": False}) 

671 if save_dot: 

672 with open(os.path.join(submit_path, "bps_generic_workflow.dot"), "w") as outfh: 

673 generic_workflow.draw(outfh, "dot") 

674 

675 _LOG.info("Starting prepare stage (creating specific implementation of workflow)") 

676 with time_this( 

677 log=_LOG, 

678 level=logging.INFO, 

679 prefix=None, 

680 msg="Prepare stage completed", 

681 mem_usage=True, 

682 mem_unit=DEFAULT_MEM_UNIT, 

683 mem_fmt=DEFAULT_MEM_FMT, 

684 ): 

685 wms_workflow = prepare(generic_workflow_config, generic_workflow, submit_path) 

686 _log_mem_usage() 

687 

688 wms_workflow_config = generic_workflow_config 

689 

690 if kwargs.get("dry_run", False): 

691 return 

692 

693 _LOG.info("Starting submit stage") 

694 with time_this( 

695 log=_LOG, 

696 level=logging.INFO, 

697 prefix=None, 

698 msg="Submit stage completed", 

699 mem_usage=True, 

700 mem_unit=DEFAULT_MEM_UNIT, 

701 mem_fmt=DEFAULT_MEM_FMT, 

702 ): 

703 submit(wms_workflow_config, wms_workflow, **kwargs) 

704 _log_mem_usage() 

705 print(f"Run Id: {wms_workflow.run_id}") 

706 print(f"Run Name: {wms_workflow.name}") 

707 

708 

709def _log_mem_usage() -> None: 

710 """Log memory usage.""" 

711 if _LOG.isEnabledFor(logging.INFO): 711 ↛ exitline 711 didn't return from function '_log_mem_usage' because the condition on line 711 was always true

712 _LOG.info( 

713 "Peak memory usage for bps process %s (main), %s (largest child process)", 

714 *tuple(f"{val.to(DEFAULT_MEM_UNIT):{DEFAULT_MEM_FMT}}" for val in get_peak_mem_usage()), 

715 ) 

716 

717 

718def batch_acquire_driver(config_file: str, **kwargs) -> None: 

719 """Create a quantum graph from pipeline definition in a batch job. 

720 

721 Parameters 

722 ---------- 

723 config_file : `str` 

724 Name of the configuration file. 

725 **kwargs 

726 Additional modifiers to the configuration. 

727 """ 

728 config = BpsConfig(config_file) 

729 translate_command_line_values(config, **kwargs) 

730 

731 found, val = config.search("saveQgraph") 

732 if found: 

733 config[".bps_defined.runQgraphFile"] = val 

734 

735 _LOG.info("Starting acquire stage (generating and/or reading quantum graph)") 

736 with time_this( 

737 log=_LOG, 

738 level=logging.INFO, 

739 prefix=None, 

740 msg="Acquire stage completed", 

741 mem_usage=True, 

742 mem_unit=DEFAULT_MEM_UNIT, 

743 mem_fmt=DEFAULT_MEM_FMT, 

744 ): 

745 _ = acquire_quantum_graph(config, out_prefix="") 

746 

747 # Copy quantum graph to staging area 

748 found, use_run_temp_space = config.search( 

749 "bpsUseRunTempSpace", opt={"curvals": {"curr_site": config[".computeSite"]}} 

750 ) 

751 if found and use_run_temp_space: 

752 found, run_temp_space = config.search( 

753 "fileDistributionEndpoint", opt={"curvals": {"curr_site": config[".computeSite"]}} 

754 ) 

755 if found: 

756 _LOG.debug("run_temp_space = %s", run_temp_space) 

757 dest = ResourcePath(run_temp_space, forceDirectory=True).join(config["qgraphFileTemplate"]) 

758 src = ResourcePath(config[".bps_defined.runQgraphFile"], forceDirectory=False) 

759 # S3 clients explicitly instantiate here to overpass this 

760 # https://stackoverflow.com/questions/52820971/is-boto3-client-thread-safe 

761 dest.exists() 

762 

763 _LOG.debug("Copying quantum graph from %s to %s", src, dest) 

764 dest.transfer_from(src, transfer="copy") 

765 else: 

766 raise KeyError("Config is missing fileDistributionEndpoint.") 

767 elif not found: 767 ↛ 770line 767 didn't jump to line 770 because the condition on line 767 was always true

768 _LOG.debug("Config is missing bpsUseRunTempSpace") 

769 

770 _log_mem_usage() 

771 

772 

773def batch_prepare_driver(config_file: str, **kwargs) -> None: 

774 """Run workflow preparation in a batch job for an existing QuantumGraph. 

775 

776 Parameters 

777 ---------- 

778 config_file : `str` 

779 Name of the configuration file. 

780 **kwargs 

781 Additional modifiers to the configuration. 

782 """ 

783 config = BpsConfig(config_file) 

784 translate_command_line_values(config, **kwargs) 

785 

786 config[".bps_defined.runQgraphFile"] = kwargs["qgraph"] 

787 submit_path = config[".bps_defined.submitPath"] 

788 

789 with time_this( 

790 log=_LOG, 

791 level=logging.INFO, 

792 prefix=None, 

793 msg="Batch preparation completed", 

794 mem_usage=True, 

795 mem_unit=DEFAULT_MEM_UNIT, 

796 mem_fmt=DEFAULT_MEM_FMT, 

797 ): 

798 batch_payload_prepare(config, prefix=submit_path) 

799 _log_mem_usage()