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 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(
|
||||||
|
|
|
||||||
|
|
@ -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,
|
||||||
|
|
|||||||
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")
|
||||||
|
|
|
||||||
|
|
@ -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,13 +41,30 @@ 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:
|
||||||
|
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}")
|
logger.info(f"Importing module: {t_module}")
|
||||||
importlib.import_module(t_module)
|
importlib.import_module(t_module)
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue
Newline not required?