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 )