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 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 queue as api_queue
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:
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
@ -188,6 +198,8 @@ def queue(
# TODO: add validation
config = parse_yaml(filepath) # type: ignore
elements = _parse_elements(element, config)
if "queue" not in config:
raise_error(f"No queue configuration found in {filepath}.")
queue_config = config.pop("queue")
kind = queue_config.pop("kind")
api_queue(

View file

@ -11,6 +11,7 @@ import typing
from pathlib import Path
from typing import Dict, List, Optional, Tuple, Union
import textwrap
import yaml
from ..datagrabber.base import BaseDataGrabber
@ -200,6 +201,22 @@ def queue(
shutil.rmtree(jobdir)
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"
logger.info(f"Writing YAML config to {str(yaml_config.absolute())}")
with open(yaml_config, "w") as f:
@ -228,6 +245,7 @@ def queue(
jobdir=jobdir,
yaml_config=yaml_config,
elements=elements, # type: ignore
config=config,
**kwargs,
)
elif kind == "SLURM":
@ -236,6 +254,7 @@ def queue(
jobdir=jobdir,
yaml_config=yaml_config,
elements=elements, # type: ignore
config=config,
**kwargs,
)
else:
@ -249,6 +268,7 @@ def _queue_condor(
jobdir: Path,
yaml_config: Path,
elements: List[Union[str, Tuple]],
config: Dict,
env: Optional[Dict[str, str]] = None,
mem: str = "8G",
cpus: int = 1,
@ -270,6 +290,8 @@ def _queue_condor(
The path to the YAML config file.
elements : list of str or tuple
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
The environment variables passed as dictionary (default None).
mem : str, optional
@ -311,8 +333,8 @@ def _queue_condor(
env_name = env["name"]
executable = "run_conda.sh"
arguments = f"{env_name} junifer"
# TODO: Copy run_conda.sh to jobdir
exec_path = jobdir / executable
logger.info(f"Copying {executable} to {str(exec_path.absolute())}")
shutil.copy(Path(__file__).parent / "res" / executable, exec_path)
make_executable(exec_path)
elif env["kind"] == "venv":
@ -351,9 +373,9 @@ def _queue_condor(
{extra_preamble}
# Logs
log = {str(log_dir.absolute())}/junifer_run_$(element).log
output = {str(log_dir.absolute())}/junifer_run_$(element).out
error = {str(log_dir.absolute())}/junifer_run_$(element).err
log = {str(log_dir.absolute())}/junifer_run_$(log_element).log
output = {str(log_dir.absolute())}/junifer_run_$(log_element).out
error = {str(log_dir.absolute())}/junifer_run_$(log_element).err
"""
submit_run_fname = jobdir / f"run_{jobname}.submit"
@ -362,7 +384,7 @@ def _queue_condor(
# Write to run submit files
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")
collect_preamble = f"""
@ -392,16 +414,22 @@ def _queue_condor(
# Now create the collect 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")
with open(dag_fname, "w") as dag_file:
# Get all subject and session names from file list
for i_job, t_elem in enumerate(elements):
str_elem = (','.join(t_elem) if isinstance(t_elem, tuple)
else t_elem)
str_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'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:
dag_file.write(f"JOB collect {submit_collect_fname}\n")
dag_file.write("PARENT ")
@ -426,6 +454,7 @@ def _queue_slurm(
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,
elements: List[Union[str, Tuple]],
config: Dict
) -> None:
"""Submit job to SLURM.
@ -440,7 +469,8 @@ def _queue_slurm(
elements : str or tuple or list[str or tuple], optional
Element(s) to process. Will be used to index the datagrabber
(default None).
config : dict
The configuration to be used for queueing the job.
"""
pass
# logger.debug("Creating SLURM job")

View file

@ -7,6 +7,9 @@
import importlib
from pathlib import Path
from typing import Dict, Union
import importlib.util
import os
import sys
import yaml
@ -38,13 +41,30 @@ def parse_yaml(filepath: Union[str, Path]) -> Dict:
# Filepath reading
with open(filepath, "r") as 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:
to_load = contents["with"]
# Convert autload modules to list
# Convert load modules to list
if not isinstance(to_load, list):
to_load = [to_load]
for t_module in to_load:
if t_module.endswith(".py"):
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)