Coverage for python/lsst/ctrl/bps/drivers.py: 61%
255 statements
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-23 10:05 +0000
« prev ^ index » next coverage.py v7.16.1, created at 2026-09-23 10:05 +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/>.
28"""Driver functions for each subcommand.
30Driver functions ensure that ensure all setup work is done before running
31the subcommand method.
32"""
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]
51import logging
52import os
53from pathlib import Path
54from typing import Any
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
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
92_LOG = logging.getLogger(__name__)
95def _init_submission_driver(config_file: str, **kwargs) -> BpsConfig:
96 """Initialize runtime environment.
98 Parameters
99 ----------
100 config_file : `str`
101 Name of the configuration file.
102 **kwargs
103 Additional modifiers to the configuration.
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()
124 submit_path = config[".bps_defined.submitPath"]
125 print(f"Submit dir: {submit_path}")
126 return config
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.
132 Parameters
133 ----------
134 config_file : `str`
135 Name of the configuration file.
136 **kwargs : `~typing.Any`
137 Additional modifiers to the configuration.
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)
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)
162 _log_mem_usage()
164 config[".bps_defined.runQgraphFile"] = qgraph_file
165 return config, qgraph
168def cluster_qgraph_driver(config_file: str, **kwargs: Any) -> tuple[BpsConfig, ClusteredQuantumGraph]:
169 """Group quanta into clusters.
171 Parameters
172 ----------
173 config_file : `str`
174 Name of the configuration file.
175 **kwargs : `~typing.Any`
176 Additional modifiers to the configuration.
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)
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()
200 _LOG.info("ClusteredQuantumGraph contains %d cluster(s)", len(clustered_qgraph))
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
212def transform_driver(config_file: str, **kwargs: Any) -> tuple[BpsConfig, GenericWorkflow]:
213 """Create a workflow for a specific workflow management system.
215 Parameters
216 ----------
217 config_file : `str`
218 Name of the configuration file.
219 **kwargs : `~typing.Any`
220 Additional modifiers to the configuration.
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"]
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()
247 num_jobs = sum(generic_workflow.job_counts.values())
248 _LOG.info("GenericWorkflow contains %d job(s) (including final)", num_jobs)
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
261def prepare_driver(config_file, **kwargs):
262 """Create a representation of the generic workflow.
264 Parameters
265 ----------
266 config_file : `str`
267 Name of the configuration file.
268 **kwargs : `~typing.Any`
269 Additional modifiers to the configuration.
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"]
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()
296 wms_workflow_config = generic_workflow_config
297 return wms_workflow_config, wms_workflow
300def submit_driver(config_file, **kwargs):
301 """Submit workflow for execution.
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)
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 )
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)
325 wms_service_class = config["wmsServiceClass"]
326 search_opts = config.get_search_opts()
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.
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]
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]
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)
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.")
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)
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()
392 _make_id_link(wms_workflow_config, wms_workflow.run_id)
394 print(f"Run Id: {wms_workflow.run_id}")
395 print(f"Run Name: {wms_workflow_config['uniqProcName']}")
398def restart_driver(wms_service, run_id):
399 """Restart a failed workflow.
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"]
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)
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")
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.
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.
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.
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"])
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
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 )
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 )
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.
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.
531 Only applicable in the context of a WMS using distributed job queues
532 (e.g., HTCondor).
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"])
543 state, message = status(
544 wms_service,
545 run_id=run_id,
546 hist=hist_days,
547 is_global=is_global,
548 )
550 _LOG.info("status: %s", state.name)
551 if message:
552 _LOG.warning(message)
554 return state.value
557def cancel_driver(wms_service, run_id, user, require_bps, pass_thru, is_global=False):
558 """Cancel submitted workflows.
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.
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)
586def ping_driver(wms_service=None, pass_thru=None):
587 """Check whether WMS services are up, reachable, and any authentication,
588 if needed, succeeds.
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.
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.
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)
611 if message:
612 if not status:
613 _LOG.info(message)
614 else:
615 _LOG.error(message)
617 # Log overall status message
618 if not status:
619 _LOG.info("Ping successful.")
620 else:
621 _LOG.error("Ping failed (%d).", status)
623 return status
626def submitcmd_driver(config_file: str, **kwargs) -> None:
627 """Submit a command for execution.
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()
650 submit_path = config[".bps_defined.submitPath"]
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()
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")
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()
688 wms_workflow_config = generic_workflow_config
690 if kwargs.get("dry_run", False):
691 return
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}")
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 )
718def batch_acquire_driver(config_file: str, **kwargs) -> None:
719 """Create a quantum graph from pipeline definition in a batch job.
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)
731 found, val = config.search("saveQgraph")
732 if found:
733 config[".bps_defined.runQgraphFile"] = val
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="")
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()
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")
770 _log_mem_usage()
773def batch_prepare_driver(config_file: str, **kwargs) -> None:
774 """Run workflow preparation in a batch job for an existing QuantumGraph.
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)
786 config[".bps_defined.runQgraphFile"] = kwargs["qgraph"]
787 submit_path = config[".bps_defined.submitPath"]
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()