Feat/import #129

Merged
fraimondo merged 5 commits from feat/import into main 2022-11-15 10:28:35 +00:00
3 changed files with 77 additions and 15 deletions

View file

@ -13,7 +13,12 @@ from typing import Dict, List, Union
import click import click
import yaml import yaml
from ..utils.logging import configure_logging, logger, warn_with_log from ..utils.logging import (
configure_logging,
logger,
warn_with_log,
raise_error,
)
from .functions import collect as api_collect from .functions import collect as api_collect
from .functions import queue as api_queue from .functions import queue as api_queue
from .functions import run as api_run from .functions import run as api_run
@ -60,6 +65,11 @@ def _parse_elements(element: str, config: Dict) -> Union[List, None]:
) )
elif elements is None: elif elements is None:
elements = config.get("elements", None) elements = config.get("elements", None)
if elements is None:
raise_error(
"The 'elements' key is set in the configuration, but its value"
" is 'None'. It is likely that there is an empty 'elements' "
"section in the yaml configuration file.")
return elements return elements
@ -188,6 +198,8 @@ def queue(
# TODO: add validation # TODO: add validation
config = parse_yaml(filepath) # type: ignore config = parse_yaml(filepath) # type: ignore
elements = _parse_elements(element, config) elements = _parse_elements(element, config)
if "queue" not in config:
raise_error(f"No queue configuration found in {filepath}.")
queue_config = config.pop("queue") queue_config = config.pop("queue")
kind = queue_config.pop("kind") kind = queue_config.pop("kind")
api_queue( api_queue(

View file

@ -11,6 +11,7 @@ import typing
from pathlib import Path from pathlib import Path
from typing import Dict, List, Optional, Tuple, Union from typing import Dict, List, Optional, Tuple, Union
import textwrap
import yaml import yaml
from ..datagrabber.base import BaseDataGrabber from ..datagrabber.base import BaseDataGrabber
@ -200,6 +201,22 @@ def queue(
shutil.rmtree(jobdir) shutil.rmtree(jobdir)
jobdir.mkdir(exist_ok=True, parents=True) jobdir.mkdir(exist_ok=True, parents=True)
if "with" in config:
to_load = config["with"]
# If there is a list of files to load, copy and remove the path
# component
fixed_load = []
if not isinstance(to_load, list):
to_load = [to_load]
for item in to_load:
if item.endswith(".py"):
logger.debug(f"Copying {item} to jobdir ({jobdir.absolute()})")
shutil.copy(item, jobdir)
fixed_load.append(Path(item).name)
else:
fixed_load.append(item)
config["with"] = fixed_load
yaml_config = jobdir / "config.yaml" yaml_config = jobdir / "config.yaml"
logger.info(f"Writing YAML config to {str(yaml_config.absolute())}") logger.info(f"Writing YAML config to {str(yaml_config.absolute())}")
with open(yaml_config, "w") as f: with open(yaml_config, "w") as f:
@ -228,6 +245,7 @@ def queue(
jobdir=jobdir, jobdir=jobdir,
yaml_config=yaml_config, yaml_config=yaml_config,
elements=elements, # type: ignore elements=elements, # type: ignore
config=config,
**kwargs, **kwargs,
) )
elif kind == "SLURM": elif kind == "SLURM":
@ -236,6 +254,7 @@ def queue(
jobdir=jobdir, jobdir=jobdir,
yaml_config=yaml_config, yaml_config=yaml_config,
elements=elements, # type: ignore elements=elements, # type: ignore
config=config,
**kwargs, **kwargs,
) )
else: else:
@ -249,6 +268,7 @@ def _queue_condor(
jobdir: Path, jobdir: Path,
yaml_config: Path, yaml_config: Path,
elements: List[Union[str, Tuple]], elements: List[Union[str, Tuple]],
config: Dict,
env: Optional[Dict[str, str]] = None, env: Optional[Dict[str, str]] = None,
mem: str = "8G", mem: str = "8G",
cpus: int = 1, cpus: int = 1,
@ -270,6 +290,8 @@ def _queue_condor(
The path to the YAML config file. The path to the YAML config file.
elements : list of str or tuple elements : list of str or tuple
Element(s) to process. Will be used to index the datagrabber. Element(s) to process. Will be used to index the datagrabber.
config : dict
The configuration to be used for queueing the job.
env : dict, optional env : dict, optional
The environment variables passed as dictionary (default None). The environment variables passed as dictionary (default None).
mem : str, optional mem : str, optional
@ -311,8 +333,8 @@ def _queue_condor(
env_name = env["name"] env_name = env["name"]
executable = "run_conda.sh" executable = "run_conda.sh"
arguments = f"{env_name} junifer" arguments = f"{env_name} junifer"
# TODO: Copy run_conda.sh to jobdir
exec_path = jobdir / executable exec_path = jobdir / executable
logger.info(f"Copying {executable} to {str(exec_path.absolute())}")
shutil.copy(Path(__file__).parent / "res" / executable, exec_path) shutil.copy(Path(__file__).parent / "res" / executable, exec_path)
make_executable(exec_path) make_executable(exec_path)
elif env["kind"] == "venv": elif env["kind"] == "venv":
@ -351,9 +373,9 @@ def _queue_condor(
{extra_preamble} {extra_preamble}
# Logs # Logs
log = {str(log_dir.absolute())}/junifer_run_$(element).log log = {str(log_dir.absolute())}/junifer_run_$(log_element).log
output = {str(log_dir.absolute())}/junifer_run_$(element).out output = {str(log_dir.absolute())}/junifer_run_$(log_element).out
error = {str(log_dir.absolute())}/junifer_run_$(element).err error = {str(log_dir.absolute())}/junifer_run_$(log_element).err
""" """
submit_run_fname = jobdir / f"run_{jobname}.submit" submit_run_fname = jobdir / f"run_{jobname}.submit"
@ -362,7 +384,7 @@ def _queue_condor(
# Write to run submit files # Write to run submit files
with open(submit_run_fname, "w") as submit_file: with open(submit_run_fname, "w") as submit_file:
submit_file.write(run_preamble) submit_file.write(textwrap.dedent(run_preamble))
submit_file.write("queue\n") submit_file.write("queue\n")
collect_preamble = f""" collect_preamble = f"""
@ -392,16 +414,22 @@ def _queue_condor(
# Now create the collect submit file # Now create the collect submit file
with open(submit_collect_fname, "w") as submit_file: with open(submit_collect_fname, "w") as submit_file:
submit_file.write(collect_preamble) # Eval preamble here submit_file.write(textwrap.dedent(collect_preamble))
submit_file.write("queue\n") submit_file.write("queue\n")
with open(dag_fname, "w") as dag_file: with open(dag_fname, "w") as dag_file:
# Get all subject and session names from file list # Get all subject and session names from file list
for i_job, t_elem in enumerate(elements): for i_job, t_elem in enumerate(elements):
str_elem = (','.join(t_elem) if isinstance(t_elem, tuple) str_elem = (
else t_elem) ",".join(t_elem) if isinstance(t_elem, tuple) else t_elem
)
log_elem = (
"_".join(t_elem) if isinstance(t_elem, tuple) else t_elem
)
dag_file.write(f"JOB run{i_job} {submit_run_fname}\n") dag_file.write(f"JOB run{i_job} {submit_run_fname}\n")
dag_file.write(f'VARS run{i_job} element="{str_elem}"\n\n') dag_file.write(
f'VARS run{i_job} element="{str_elem} '
f'log_element="{log_elem}"\n\n')
if collect is True: if collect is True:
dag_file.write(f"JOB collect {submit_collect_fname}\n") dag_file.write(f"JOB collect {submit_collect_fname}\n")
dag_file.write("PARENT ") dag_file.write("PARENT ")
@ -426,6 +454,7 @@ def _queue_slurm(
jobdir: Path, jobdir: Path,
synchon commented 2022-11-15 09:31:58 +00:00 (Migrated from github.com)

Newline not required?

Newline not required?
yaml_config: Path, yaml_config: Path,
elements: List[Union[str, Tuple]], elements: List[Union[str, Tuple]],
config: Dict
) -> None: ) -> None:
"""Submit job to SLURM. """Submit job to SLURM.
@ -440,7 +469,8 @@ def _queue_slurm(
elements : str or tuple or list[str or tuple], optional elements : str or tuple or list[str or tuple], optional
Element(s) to process. Will be used to index the datagrabber Element(s) to process. Will be used to index the datagrabber
(default None). (default None).
config : dict
The configuration to be used for queueing the job.
""" """
pass pass
# logger.debug("Creating SLURM job") # logger.debug("Creating SLURM job")

View file

@ -7,6 +7,9 @@
import importlib import importlib
from pathlib import Path from pathlib import Path
from typing import Dict, Union from typing import Dict, Union
import importlib.util
import os
import sys
import yaml import yaml
@ -38,14 +41,31 @@ def parse_yaml(filepath: Union[str, Path]) -> Dict:
# Filepath reading # Filepath reading
with open(filepath, "r") as f: with open(filepath, "r") as f:
contents = yaml.safe_load(f) contents = yaml.safe_load(f)
# Autload modules if "elements" in contents:
if contents["elements"] is None:
raise_error(
"The elements key was defined but its content is empty. "
"Please define the elements to operate on or remove the key.")
# load modules
if "with" in contents: if "with" in contents:
to_load = contents["with"] to_load = contents["with"]
# Convert autload modules to list # Convert load modules to list
if not isinstance(to_load, list): if not isinstance(to_load, list):
to_load = [to_load] to_load = [to_load]
for t_module in to_load: for t_module in to_load:
logger.info(f"Importing module: {t_module}") if t_module.endswith(".py"):
importlib.import_module(t_module) logger.debug(f"Importing file: {t_module}")
file_path = Path(os.getcwd()) / t_module
if not file_path.exists():
raise_error(
f"File in 'with' section does not exist: {file_path}")
spec = importlib.util.spec_from_file_location(
t_module, file_path)
module = importlib.util.module_from_spec(spec) # type: ignore
sys.modules[t_module] = module
spec.loader.exec_module(module) # type: ignore
else:
logger.info(f"Importing module: {t_module}")
importlib.import_module(t_module)
return contents return contents