Feat/import #129
3 changed files with 77 additions and 15 deletions
|
|
@ -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(
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
||||
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")
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue
Newline not required?