Back to home page

EIC code displayed by LXR

 
 

    


File indexing completed on 2026-08-12 08:24:57

0001 """Runtime builder helpers for config-driven AID2E execution."""
0002 
0003 from __future__ import annotations
0004 
0005 import importlib
0006 from typing import Any, Dict, Optional, Union
0007 
0008 from aid2e.utilities.configurations.optimizer_config import OptimizerConfiguration
0009 from aid2e.utilities.configurations.problem_config import ProblemConfiguration
0010 from aid2e.utilities.configurations.scheduler_config import SchedulerConfiguration
0011 from aid2e.utilities.configurations.scheduler_cascade import resolve_scheduler_cascade
0012 from aid2e.utilities.configurations.workflow_config import (
0013     JobDefinition,
0014     WorkflowDefinition,
0015     WorkflowsConfiguration,
0016 )
0017 
0018 
0019 def infer_optimizer_backend(optimizer_cfg: OptimizerConfiguration) -> str:
0020     """Infer optimizer backend name from optimizer configuration."""
0021     name = (optimizer_cfg.name or "").strip().lower()
0022     opt_type = (optimizer_cfg.type or "").strip().lower()
0023     params = optimizer_cfg.parameters or {}
0024     algorithm = str(params.get("algorithm", "")).strip().lower()
0025 
0026     tokens = {name, opt_type, algorithm}
0027     if tokens & {"ax", "bo", "mobo", "bayesian"}:
0028         return "ax"
0029     if tokens & {"pymoo", "ga", "nsga2", "nsga3", "moead", "evolutionary"}:
0030         return "pymoo"
0031 
0032     raise ValueError(
0033         "Could not infer optimizer backend from optimizer config. "
0034         "Use optimizer.name/type like 'ax' or 'pymoo'."
0035     )
0036 
0037 
0038 def build_optimizer_from_config(
0039     problem_cfg: ProblemConfiguration,
0040     optimizer_cfg: OptimizerConfiguration,
0041     *,
0042     backend: Optional[str] = None,
0043 ):
0044     """Build a concrete optimizer from parsed configs."""
0045     backend_name = (backend or infer_optimizer_backend(optimizer_cfg)).lower()
0046     objective_names = [obj.name for obj in problem_cfg.objectives]
0047     objective_directions = {obj.name: obj.direction for obj in problem_cfg.objectives}
0048     params = dict(optimizer_cfg.parameters or {})
0049 
0050     if backend_name == "ax":
0051         from aid2e.optimizers.ax import AxOptimizer, AxOptimizerConfig
0052 
0053         ax_cfg = AxOptimizerConfig(**params)
0054         return AxOptimizer(
0055             search_space=problem_cfg.design_config,
0056             config=ax_cfg,
0057             objective_names=objective_names,
0058             objective_directions=objective_directions,
0059             seed=ax_cfg.seed,
0060         )
0061 
0062     if backend_name == "pymoo":
0063         from aid2e.optimizers.pymoo import PyMOOOptimizer, PyMOOOptimizerConfig
0064 
0065         if "seed" not in params:
0066             params["seed"] = None
0067         pymoo_cfg = PyMOOOptimizerConfig(**params)
0068         return PyMOOOptimizer(
0069             search_space=problem_cfg.design_config,
0070             config=pymoo_cfg,
0071             objective_names=objective_names,
0072             seed=pymoo_cfg.seed,
0073         )
0074 
0075     raise ValueError(f"Unsupported optimizer backend: {backend_name}")
0076 
0077 
0078 def build_scheduler_runtime_config(
0079     scheduler_cfg: Optional[SchedulerConfiguration],
0080 ) -> Optional[Dict[str, Any]]:
0081     """Convert SchedulerConfiguration to DAGExecutor scheduler payload."""
0082     if scheduler_cfg is None:
0083         return None
0084 
0085     runner_type = scheduler_cfg.runner_type
0086     params = dict(scheduler_cfg.parameters or {})
0087 
0088     if not params:
0089         payload = scheduler_cfg.model_dump()
0090         if runner_type == "JobLibRunner":
0091             params = dict(payload.get("joblib") or {})
0092         elif runner_type == "PanDAiDDSRunner":
0093             params = dict(payload.get("pandaidds") or payload.get("panda") or {})
0094 
0095     if runner_type == "JobLibRunner":
0096         from aid2e.schedulers.JobLib import JobLibRunnerConfig
0097 
0098         cfg = JobLibRunnerConfig(**params)
0099     elif runner_type == "SlurmRunner":
0100         from aid2e.schedulers.Slurm import SlurmRunnerConfig
0101 
0102         cfg = SlurmRunnerConfig(**params)
0103     elif runner_type == "PanDAiDDSRunner":
0104         from aid2e.schedulers.PanDAiDDS import PanDAiDDSRunnerConfig
0105 
0106         cfg = PanDAiDDSRunnerConfig(**params)
0107     else:
0108         raise ValueError(f"Unsupported scheduler runner_type: {runner_type}")
0109 
0110     return {"runner_type": runner_type, "config": cfg}
0111 
0112 
0113 def build_scheduler_from_config(scheduler_cfg: Optional[SchedulerConfiguration]):
0114     """Build a scheduler instance from scheduler configuration."""
0115     runtime_cfg = build_scheduler_runtime_config(scheduler_cfg)
0116     if runtime_cfg is None:
0117         return None
0118 
0119     runner_type = runtime_cfg["runner_type"]
0120     cfg_obj = runtime_cfg["config"]
0121     if runner_type == "JobLibRunner":
0122         from aid2e.schedulers.JobLib import JobLibScheduler
0123 
0124         return JobLibScheduler(config=cfg_obj)
0125     if runner_type == "SlurmRunner":
0126         from aid2e.schedulers.Slurm import SlurmScheduler
0127 
0128         return SlurmScheduler(config=cfg_obj)
0129     if runner_type == "PanDAiDDSRunner":
0130         from aid2e.schedulers.PanDAiDDS import PanDAiDDSScheduler
0131 
0132         return PanDAiDDSScheduler(config=cfg_obj)
0133 
0134     raise ValueError(f"Unsupported scheduler runner_type: {runner_type}")
0135 
0136 
0137 def _resolve_callable(spec: str):
0138     """Resolve a callable from '<module>:<symbol>' or '<module>.<symbol>'."""
0139     if ":" in spec:
0140         module_name, symbol_name = spec.split(":", 1)
0141     else:
0142         module_name, symbol_name = spec.rsplit(".", 1)
0143     module = importlib.import_module(module_name)
0144     return getattr(module, symbol_name)
0145 
0146 
0147 def _resolve_workflow_python_callables(
0148     workflow: WorkflowDefinition,
0149 ) -> WorkflowDefinition:
0150     """Resolve string payload callable references for workflow jobs."""
0151     wf = workflow.model_copy(deep=True)
0152     for branch in wf.branches:
0153         for stage in branch.stages:
0154             resolved_jobs = []
0155             for job in stage.jobs:
0156                 payload = dict(job.payload or {})
0157                 callable_spec = payload.get("python_callable")
0158                 if isinstance(callable_spec, str):
0159                     payload["python_callable"] = _resolve_callable(callable_spec)
0160                 resolved_jobs.append(
0161                     JobDefinition(
0162                         name=job.name,
0163                         command=job.command,
0164                         payload=payload,
0165                         rule=job.rule,
0166                         resources=job.resources,
0167                         outputs=job.outputs,
0168                     )
0169                 )
0170             stage.jobs = resolved_jobs
0171     return wf
0172 
0173 
0174 def select_workflow(
0175     workflows_cfg: Union[WorkflowDefinition, WorkflowsConfiguration],
0176     workflow_name: Optional[str] = None,
0177 ) -> WorkflowDefinition:
0178     """Select one workflow definition from a workflow config container."""
0179     if isinstance(workflows_cfg, WorkflowDefinition):
0180         return workflows_cfg
0181 
0182     workflows = workflows_cfg.workflows
0183     if not workflows:
0184         raise ValueError("No workflows found in workflow configuration")
0185 
0186     if workflow_name:
0187         for wf in workflows:
0188             if wf.name == workflow_name:
0189                 return wf
0190         raise ValueError(f"Workflow '{workflow_name}' not found in config")
0191 
0192     return workflows[0]
0193 
0194 
0195 def build_workflow_executor_from_config(
0196     workflows_cfg: Union[WorkflowDefinition, WorkflowsConfiguration],
0197     *,
0198     problem_cfg: Optional[ProblemConfiguration] = None,
0199     scheduler_cfg: Optional[SchedulerConfiguration] = None,
0200     workflow_name: Optional[str] = None,
0201     base_output_dir: str = "/tmp/aid2e_runs",
0202     log_level: str = "INFO",
0203 ):
0204     """Build a DAGExecutor from workflow + scheduler configuration."""
0205     from aid2e.utilities.workflows import DAGExecutor
0206 
0207     workflow = select_workflow(workflows_cfg, workflow_name=workflow_name)
0208     if problem_cfg is not None and workflow.objectives:
0209         raise ValueError(
0210             "Canonical full-config workflows must not repeat 'workflows[].objectives'. "
0211             "Use 'problem.objectives' as the single source of truth."
0212         )
0213     resolved = _resolve_workflow_python_callables(workflow)
0214     if problem_cfg is not None:
0215         resolved.objectives = list(problem_cfg.objectives)
0216 
0217     workflow_global_scheduler = (
0218         workflows_cfg.global_scheduler
0219         if isinstance(workflows_cfg, WorkflowsConfiguration)
0220         else None
0221     )
0222     global_scheduler = (
0223         scheduler_cfg if scheduler_cfg is not None else workflow_global_scheduler
0224     )
0225 
0226     def resolve_stage_scheduler(branch, stage):
0227         return build_scheduler_runtime_config(
0228             resolve_scheduler_cascade(
0229                 stage_scheduler=stage.scheduler,
0230                 branch_scheduler=branch.scheduler,
0231                 workflow_scheduler=resolved.scheduler,
0232                 global_scheduler=global_scheduler,
0233             )
0234         )
0235 
0236     runtime_scheduler_cfg = build_scheduler_runtime_config(global_scheduler)
0237 
0238     return DAGExecutor(
0239         workflow=resolved,
0240         base_output_dir=base_output_dir,
0241         log_level=log_level,
0242         problem_config=problem_cfg,
0243         scheduler_config=runtime_scheduler_cfg,
0244         scheduler_config_resolver=resolve_stage_scheduler,
0245     )