From 9fe0939813198e0912a19cdd338012cbc1b08341 Mon Sep 17 00:00:00 2001 From: doxav Date: Mon, 22 Sep 2025 08:13:27 +0200 Subject: [PATCH 1/6] initial working version of GEPA bench test --- opto/trainer/algorithms/gepa_algorithms.py | 652 ++++++++++++++++++ .../test_gepa_benchmark.py | 94 +++ 2 files changed, 746 insertions(+) create mode 100644 opto/trainer/algorithms/gepa_algorithms.py create mode 100644 tests/llm_optimizers_tests/test_gepa_benchmark.py diff --git a/opto/trainer/algorithms/gepa_algorithms.py b/opto/trainer/algorithms/gepa_algorithms.py new file mode 100644 index 00000000..588cdbad --- /dev/null +++ b/opto/trainer/algorithms/gepa_algorithms.py @@ -0,0 +1,652 @@ +# opto/trainer/algorithms/gepa_algorithms.py +# GEPA (+Merge) algorithms for Trace +# - GEPAUCBSearch: subclass of UCBSearchAlgorithm +# - GEPABeamPareto: subclass of BeamsearchAlgorithm (Pareto select + single-parent incremental) +# - GEPATrainer: subclass of Trainer (minimal GEPA loop) +# +# All default to OptoPrimeV2 if optimizer=None. + +from __future__ import annotations +import copy +import math +import random +from dataclasses import dataclass, field +from typing import Any, Dict, List, Optional, Tuple + +import numpy as np + +from opto.optimizers.optoprime_v2 import OptoPrimeV2 +from opto.trace.nodes import ParameterNode +from opto.trainer.algorithms.UCBsearch import UCBSearchAlgorithm +from opto.trainer.algorithms.beamsearch_algorithm import BeamsearchAlgorithm +from opto.trainer.algorithms.algorithm import Trainer +from opto.trainer.algorithms.basic_algorithms import ( + evaluate, + batchify, + standard_optimization_step, +) +from opto.trainer.utils import async_run +from opto.optimizers.utils import print_color + + +# ----------------------------- Utilities ----------------------------- # + +@dataclass +class Candidate: + params: Dict[ParameterNode, Any] + eval_vector: List[float] # per-instance scores on fixed Pareto subset + mean: float + id: int + parent_ids: Tuple[int, ...] = field(default_factory=tuple) + ancestors: set = field(default_factory=set) + created_iter: int = 0 + wins: int = 0 # updated by Pareto accounting + meta: Dict[str, Any] = field(default_factory=dict) # freeform + +def _eval_on_subset(agent, guide, xs, infos, *, num_threads: Optional[int], desc: str) -> List[float]: + return evaluate(agent, guide, xs, infos, min_score=None, num_threads=num_threads, description=desc) + +def _compute_pareto_counts(cands: List[Candidate]) -> None: + """ + "Best-for-at-least-one-instance" winners. + For each position m in eval vectors, find argmax candidate and credit a win. + """ + if not cands: + return + L = len(cands[0].eval_vector) + # Reset + for c in cands: + c.wins = 0 + # Credit wins + for m in range(L): + best_idx = None + best_val = -float("inf") + for i, c in enumerate(cands): + v = c.eval_vector[m] if m < len(c.eval_vector) else -float("inf") + if v > best_val: + best_val, best_idx = v, i + if best_idx is not None: + cands[best_idx].wins += 1 + +def _pareto_sample(cands: List[Candidate], *, temperature: float = 1.0, rng: random.Random) -> Candidate: + """ + Sample a parent from union of per-instance winners, proportional to wins^1/T. + """ + if not cands: + raise ValueError("Empty candidate buffer.") + _compute_pareto_counts(cands) + wins = np.array([max(1, c.wins) for c in cands], dtype=float) # avoid zero + if temperature <= 0: + # Deterministic pick + return cands[int(wins.argmax())] + weights = wins ** (1.0 / max(1e-6, temperature)) + probs = weights / (weights.sum() if weights.sum() > 0 else 1.0) + idx = rng.choices(range(len(cands)), weights=probs, k=1)[0] + return cands[idx] + +def _uniform_merge_params(a: Dict[ParameterNode, Any], b: Dict[ParameterNode, Any], rng: random.Random) -> Dict[ParameterNode, Any]: + """ + Simple, robust "crossover": per-parameter uniform pick between parents. + (System-aware enough for prompt/code params, cheap, and safe.) + """ + keys = set(a.keys()) | set(b.keys()) + merged: Dict[ParameterNode, Any] = {} + for p in keys: + if p in a and p in b: + merged[p] = copy.deepcopy(a[p] if rng.random() < 0.5 else b[p]) + elif p in a: + merged[p] = copy.deepcopy(a[p]) + else: + merged[p] = copy.deepcopy(b[p]) + return merged + +def _maybe_merge(buffer: List[Candidate], + *, + agent, + guide, + pareto_inputs: List[Any], + pareto_infos: List[Any], + num_threads: Optional[int], + rng: random.Random, + tried_pairs: set, + max_tries: int = 8) -> Optional[Candidate]: + """ + Try merging two non-lineage candidates once; return merged if better than both parents' mean, else None. + """ + if len(buffer) < 2: + return None + # Prefer winners + _compute_pareto_counts(buffer) + pool = sorted(buffer, key=lambda c: (c.wins, c.mean), reverse=True) + + # Try a few distinct pairs + for _ in range(max_tries): + i, j = rng.sample(range(len(pool)), 2) + a, b = pool[i], pool[j] + if a.id == b.id: + continue + if a.id in b.ancestors or b.id in a.ancestors: + continue # avoid direct ancestry + key = tuple(sorted((a.id, b.id))) + if key in tried_pairs: + continue + tried_pairs.add(key) + + merged_params = _uniform_merge_params(a.params, b.params, rng) + # Evaluate merged on Pareto subset + original_params = {p: copy.deepcopy(p.data) for p in agent.parameters()} + try: + # load params to agent + from opto.optimizers.optimizer import Optimizer # type: ignore + # We only need the parameters dict projection; we can set via optimizer.update if available + # But we don't have an optimizer here; use ParameterNode._set + for p, v in merged_params.items(): + p._set(v) + + vec = _eval_on_subset(agent, guide, pareto_inputs, pareto_infos, num_threads=num_threads, + desc="GEPA+Merge: evaluating merged") + mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + finally: + # restore original + for p, v in original_params.items(): + p._set(v) + + if mean > max(a.mean, b.mean): + merged = Candidate(params=merged_params, + eval_vector=vec, + mean=mean, + id=-1, # to be set by caller + parent_ids=(a.id, b.id), + ancestors=set(a.ancestors) | set(b.ancestors) | {a.id, b.id}, + created_iter=0) + return merged + return None + + +def _ensure_optimizer(agent, optimizer): + if optimizer is not None: + return optimizer + params = [p for p in agent.parameters()] # List[ParameterNode] + return OptoPrimeV2(parameters=params) + + +def _train_step_generate_child(agent, guide, optimizer, train_xs, train_infos, *, verbose=False, num_threads=None): + """ + Single-parent, incremental evolution "mutation": run forward on a minibatch to get batched feedback, + then optimizer.step(bypassing=True) to obtain a new candidate param dict (without applying). + """ + use_async = num_threads is not None and num_threads > 1 + if use_async: + outputs = async_run([lambda a,x,g,info: standard_optimization_step(a, x, g, info)] * len(train_xs), + args_list=[(agent, x, guide, info) for x, info in zip(train_xs, train_infos)], + max_workers=num_threads, + description="GEPA forward (mutate parent)") + # outputs: List[(target, score, feedback)] + else: + outputs = [standard_optimization_step(agent, x, guide, info) for x, info in zip(train_xs, train_infos)] + + scores, targets, feedbacks = [], [], [] + for target, score, feedback in outputs: + scores.append(score) + targets.append(target) + feedbacks.append(feedback) + + target_batch = batchify(*targets) + feedback_batch = batchify(*feedbacks).data + + optimizer.zero_feedback() + optimizer.backward(target_batch, feedback_batch) + try: + update_dict = optimizer.step(bypassing=True, verbose=("output" if verbose else False)) + if not isinstance(update_dict, dict) or len(update_dict) == 0: + # Fallback: treat current as child (rare) + update_dict = {p: copy.deepcopy(p.data) for p in optimizer.parameters} + except Exception as e: + print_color(f"[GEPA] optimizer.step error: {e}", "red") + update_dict = {} + return update_dict, (None if not scores or any(s is None for s in scores) else float(np.mean(scores))) + + +def _apply_params(optimizer, param_dict: Dict[ParameterNode, Any]): + """Load param dict into the agent via optimizer.update (preserves projections).""" + optimizer.update(param_dict) + + +# ======================= Variant 1: GEPA + Merge (UCB subclass) ======================= # + +class GEPAUCBSearch(UCBSearchAlgorithm): + """ + GEPA (+Merge) implemented atop UCBSearchAlgorithm. + Differences vs base UCB: + - Fixed Pareto subset (D_pareto) and per-instance vectors kept for each candidate + - Parent selection = Pareto "best-for-at-least-one" sampling (wins-weighted); UCB used only for eviction fallback + - Single-parent incremental mutation via a minibatch + - Optional periodic Merge crossover (uniform per-parameter) with desirability checks + """ + + def __init__(self, + agent, + optimizer=None, + *, + max_buffer_size: int = 16, + ucb_exploration_factor: float = 0.8, + rng_seed: int = 7, + logger=None, + num_threads: Optional[int] = None): + optimizer = _ensure_optimizer(agent, optimizer) + super().__init__(agent, optimizer, + max_buffer_size=max_buffer_size, + ucb_exploration_factor=ucb_exploration_factor, + logger=logger, + num_threads=num_threads) + self.rng = random.Random(rng_seed) + self._pareto_inputs: List[Any] = [] + self._pareto_infos: List[Any] = [] + self._id_counter = 0 + + def _next_id(self) -> int: + self._id_counter += 1 + return self._id_counter + + def _evaluate_on_pareto(self, params_dict: Dict[ParameterNode, Any], guide, *, num_threads) -> Tuple[List[float], float]: + original_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + try: + _apply_params(self.optimizer, params_dict) + vec = _eval_on_subset(self.agent, guide, self._pareto_inputs, self._pareto_infos, + num_threads=num_threads, desc="GEPA: evaluate on Pareto subset") + mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + return vec, mean + finally: + _apply_params(self.optimizer, original_params) + + def _select_pareto_parent(self, cand_buffer: List[Candidate]) -> Candidate: + return _pareto_sample(cand_buffer, temperature=1.0, rng=self.rng) + + def train(self, + guide, + train_dataset: Dict[str, List[Any]], + *, + validate_dataset: Optional[Dict[str, List[Any]]] = None, + pareto_subset_size: int = 24, + num_search_iterations: int = 120, + train_batch_size: int = 2, + merge_every: int = 6, + log_frequency: Optional[int] = None, + save_frequency: Optional[int] = None, + save_path: str = "checkpoints/gepa_ucb_agent.pkl", + verbose: bool = False, + num_threads: Optional[int] = None) -> Tuple[Dict[str, Any], float]: + """ + GEPA search loop with Pareto sampling + (optional) Merge. + """ + num_threads = num_threads or self.num_threads + log_frequency = log_frequency or 5 + validate_ds = validate_dataset or train_dataset + + # Fix a Pareto subset (small, stable) to compute per-instance vectors + assert len(validate_ds["inputs"]) > 0, "Empty dataset." + idxs = np.random.choice(len(validate_ds["inputs"]), + min(pareto_subset_size, len(validate_ds["inputs"])), + replace=False) + self._pareto_inputs = [validate_ds["inputs"][i] for i in idxs] + self._pareto_infos = [validate_ds["infos"][i] for i in idxs] + + buffer: List[Candidate] = [] + tried_merges: set = set() + + # Seed with current params + base_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + v0, m0 = self._evaluate_on_pareto(base_params, guide, num_threads=num_threads) + buffer.append(Candidate(params=base_params, eval_vector=v0, mean=m0, id=self._next_id(), ancestors=set())) + print_color(f"[GEPA] Seed candidate mean={m0:.4f}", "cyan") + + metrics = {"best_means": [], "new_child_means": [], "merge_accepts": 0, "total_merges": 0} + + for it in range(1, num_search_iterations + 1): + # Select parent by Pareto winners + parent = self._select_pareto_parent(buffer) + _apply_params(self.optimizer, parent.params) + + # Sample train minibatch + train_size = min(train_batch_size, len(train_dataset["inputs"])) + tr_idxs = np.random.choice(len(train_dataset["inputs"]), train_size, replace=False) + train_xs = [train_dataset["inputs"][i] for i in tr_idxs] + train_info = [train_dataset["infos"][i] for i in tr_idxs] + + # Generate child via one incremental step + update_dict, train_batch_mean = _train_step_generate_child( + self.agent, guide, self.optimizer, train_xs, train_info, verbose=verbose, num_threads=num_threads + ) + if not update_dict: + print_color("[GEPA] Empty child update; skipping.", "yellow") + continue + + # Evaluate child on Pareto subset + child_vec, child_mean = self._evaluate_on_pareto(update_dict, guide, num_threads=num_threads) + child = Candidate(params=update_dict, + eval_vector=child_vec, + mean=child_mean, + id=self._next_id(), + parent_ids=(parent.id,), + ancestors=set(parent.ancestors) | {parent.id}, + created_iter=it) + buffer.append(child) + metrics["new_child_means"].append(child_mean) + print_color(f"[GEPA] iter {it}: child mean={child_mean:.4f} (train-batch≈{train_batch_mean})", "green") + + # Optional Merge + if merge_every and (it % merge_every == 0): + metrics["total_merges"] += 1 + merged = _maybe_merge(buffer, + agent=self.agent, guide=guide, + pareto_inputs=self._pareto_inputs, + pareto_infos=self._pareto_infos, + num_threads=num_threads, + rng=self.rng, + tried_pairs=tried_merges) + if merged is not None: + merged.id = self._next_id() + merged.created_iter = it + buffer.append(merged) + metrics["merge_accepts"] += 1 + print_color(f"[GEPA] Merge accepted: mean={merged.mean:.4f}", "magenta") + + # Keep buffer bounded: remove the candidate with lowest (wins, mean) + if len(buffer) > self.max_buffer_size: + _compute_pareto_counts(buffer) + buffer.sort(key=lambda c: (c.wins, c.mean)) + evicted = buffer.pop(0) + print_color(f"[GEPA] Evicted cand#{evicted.id} (wins={evicted.wins}, mean={evicted.mean:.4f})", "yellow") + + # Track & log + best = max(buffer, key=lambda c: c.mean) + metrics["best_means"].append(best.mean) + if it % log_frequency == 0: + self.logger.log("GEPA best mean", best.mean, it, color="green") + + # Save best candidate snapshot (optional) + if save_frequency and it % save_frequency == 0: + _apply_params(self.optimizer, best.params) + self.save_agent(save_path, it) + + # Load best into the agent and return + best = max(buffer, key=lambda c: c.mean) if buffer else buffer[0] + _apply_params(self.optimizer, best.params) + return metrics, float(best.mean) + + +# ================= Variant 2: Beamsearch subclass with Pareto select ================= # + +class GEPABeamPareto(BeamsearchAlgorithm): + """ + BeamsearchAlgorithm retrofit: + - override select() to a Pareto "best-for-at-least-one" selector + - replace deep beam expansion with GEPA’s single-parent incremental evolution + """ + + def __init__(self, + agent, + optimizer=None, + *, + rng_seed: int = 11, + logger=None, + num_threads: Optional[int] = None): + optimizer = _ensure_optimizer(agent, optimizer) + super().__init__(agent, optimizer, num_threads=num_threads, logger=logger) + self.rng = random.Random(rng_seed) + + # We keep a Pareto select helper that returns (selected_params, wins, scores) + def select(self, + candidates: List[Dict[ParameterNode, Any]], + validate_guide, + validation_mini_dataset, + beam_width: int, + num_threads: int = None, + min_score: float = None, + return_scores: bool = False): + """ + Override to Pareto union-of-winners on the mini validation batch. + """ + # Evaluate each candidate to a vector on the mini validation + cand_objs: List[Candidate] = [] + current_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + try: + for idx, params in enumerate(candidates): + _apply_params(self.optimizer, params) + vec = evaluate(self.agent, + validate_guide, + validation_mini_dataset['inputs'], + validation_mini_dataset['infos'], + min_score=min_score, + num_threads=num_threads, + description=f"Validating candidate {idx+1}/{len(candidates)} (Pareto)") + mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + cand_objs.append(Candidate(params=params, eval_vector=vec, mean=mean, id=idx)) + finally: + _apply_params(self.optimizer, current_params) + + # Compute wins and select top "beam_width" by (wins, mean) + _compute_pareto_counts(cand_objs) + cand_objs.sort(key=lambda c: (c.wins, c.mean), reverse=True) + selected = cand_objs[: min(beam_width, len(cand_objs))] + sel_params = [c.params for c in selected] + sel_scores = [c.mean for c in selected] + if return_scores: + return sel_params, sel_scores + return sel_params + + # Replace beam "train" with GEPA-style incremental loop (keeps BeamsearchAlgorithm API) + def train(self, + guide, + train_dataset, + *, + validate_dataset=None, + pareto_subset_size: int = 24, + num_search_iterations: int = 120, + train_batch_size: int = 2, + merge_every: int = 6, + log_frequency: Optional[int] = None, + save_frequency: Optional[int] = None, + save_path: str = "checkpoints/gepa_beam_agent.pkl", + verbose: bool = False, + num_threads: Optional[int] = None): + num_threads = num_threads or self.num_threads + log_frequency = log_frequency or 5 + validate_ds = validate_dataset or train_dataset + + # Fix Pareto subset for this run + idxs = np.random.choice(len(validate_ds["inputs"]), + min(pareto_subset_size, len(validate_ds["inputs"])), + replace=False) + pareto_inputs = [validate_ds["inputs"][i] for i in idxs] + pareto_infos = [validate_ds["infos"][i] for i in idxs] + + # Seed buffer + buffer: List[Candidate] = [] + base_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + # Evaluate seed + current_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + try: + _apply_params(self.optimizer, base_params) + vec = evaluate(self.agent, guide, pareto_inputs, pareto_infos, + min_score=None, num_threads=num_threads, + description="GEPA(beam): seed evaluation") + finally: + _apply_params(self.optimizer, current_params) + m0 = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + buffer.append(Candidate(params=base_params, eval_vector=vec, mean=m0, id=0, ancestors=set())) + tried_merges: set = set() + + best_mean = m0 + for it in range(1, num_search_iterations + 1): + # Pareto-select parent and mutate + _compute_pareto_counts(buffer) + parent = _pareto_sample(buffer, temperature=1.0, rng=self.rng) + _apply_params(self.optimizer, parent.params) + + # Make a child + k = min(train_batch_size, len(train_dataset["inputs"])) + tr = np.random.choice(len(train_dataset["inputs"]), k, replace=False) + train_xs = [train_dataset["inputs"][i] for i in tr] + train_in = [train_dataset["infos"][i] for i in tr] + + update_dict, _ = _train_step_generate_child(self.agent, guide, self.optimizer, train_xs, train_in, + verbose=verbose, num_threads=num_threads) + if not update_dict: + continue + + # Evaluate child on Pareto subset + current_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + try: + _apply_params(self.optimizer, update_dict) + vec = evaluate(self.agent, guide, pareto_inputs, pareto_infos, min_score=None, + num_threads=num_threads, description="GEPA(beam): child eval") + finally: + _apply_params(self.optimizer, current_params) + mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + buffer.append(Candidate(params=update_dict, eval_vector=vec, mean=mean, id=len(buffer), + parent_ids=(parent.id,), ancestors=set(parent.ancestors) | {parent.id})) + best_mean = max(best_mean, mean) + if it % log_frequency == 0: + self.logger.log("GEPA(beam) best mean", best_mean, it, color="green") + + # Periodic merge + if merge_every and it % merge_every == 0: + merged = _maybe_merge(buffer, + agent=self.agent, guide=guide, + pareto_inputs=pareto_inputs, pareto_infos=pareto_infos, + num_threads=num_threads, rng=self.rng, tried_pairs=tried_merges) + if merged is not None: + merged.id = len(buffer) + buffer.append(merged) + + # Trim buffer softly (keep top by (wins, mean)) + if len(buffer) > 16: + _compute_pareto_counts(buffer) + buffer.sort(key=lambda c: (c.wins, c.mean), reverse=True) + buffer[:] = buffer[:16] + + # Optional save + if save_frequency and it % save_frequency == 0: + best = max(buffer, key=lambda c: c.mean) + _apply_params(self.optimizer, best.params) + self.save_agent(save_path, it) + + best = max(buffer, key=lambda c: c.mean) + _apply_params(self.optimizer, best.params) + return {"best_mean": best.mean}, float(best.mean) + + +# =================== Variant 3: Minimal GEPA on AlgorithmBase =================== # + +class GEPAAlgorithmBase(Trainer): + """ + Lightweight GEPA (+Merge) with only Trainer dependency. + Useful when you want the simplest control loop with your own logging/saving. + """ + + def __init__(self, + agent, + optimizer=None, + *, + rng_seed: int = 13, + logger=None, + num_threads: Optional[int] = None): + super().__init__(agent, num_threads=num_threads, logger=logger) + self.optimizer = _ensure_optimizer(agent, optimizer) + self.rng = random.Random(rng_seed) + + def train(self, + guide, + train_dataset, + *, + validate_dataset=None, + pareto_subset_size: int = 24, + num_iters: int = 100, + train_batch_size: int = 2, + merge_every: int = 5, + num_threads: Optional[int] = None, + save_path: Optional[str] = None): + num_threads = num_threads or self.num_threads + validate_ds = validate_dataset or train_dataset + + # Pareto subset + idxs = np.random.choice(len(validate_ds["inputs"]), + min(pareto_subset_size, len(validate_ds["inputs"])), + replace=False) + xsP = [validate_ds["inputs"][i] for i in idxs] + isP = [validate_ds["infos"][i] for i in idxs] + + # Seed + buffer: List[Candidate] = [] + base_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + original = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + try: + _apply_params(self.optimizer, base_params) + vec = evaluate(self.agent, guide, xsP, isP, min_score=None, num_threads=num_threads, + description="GEPA(base): seed eval") + finally: + _apply_params(self.optimizer, original) + m0 = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + buffer.append(Candidate(params=base_params, eval_vector=vec, mean=m0, id=0, ancestors=set())) + tried_merges: set = set() + + for it in range(1, num_iters + 1): + # Parent select + _compute_pareto_counts(buffer) + parent = _pareto_sample(buffer, temperature=1.0, rng=self.rng) + _apply_params(self.optimizer, parent.params) + + # Child + k = min(train_batch_size, len(train_dataset["inputs"])) + tr = np.random.choice(len(train_dataset["inputs"]), k, replace=False) + tx = [train_dataset["inputs"][i] for i in tr] + ti = [train_dataset["infos"][i] for i in tr] + update_dict, _ = _train_step_generate_child(self.agent, guide, self.optimizer, tx, ti, + verbose=False, num_threads=num_threads) + if not update_dict: + continue + + # Eval child + original = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + try: + _apply_params(self.optimizer, update_dict) + vec = evaluate(self.agent, guide, xsP, isP, min_score=None, num_threads=num_threads, + description="GEPA(base): child eval") + finally: + _apply_params(self.optimizer, original) + mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + buffer.append(Candidate(params=update_dict, eval_vector=vec, mean=mean, id=len(buffer), + parent_ids=(parent.id,), ancestors=set(parent.ancestors) | {parent.id})) + + # Merge + if merge_every and it % merge_every == 0: + merged = _maybe_merge(buffer, + agent=self.agent, guide=guide, + pareto_inputs=xsP, pareto_infos=isP, + num_threads=num_threads, rng=self.rng, tried_pairs=tried_merges) + if merged is not None: + merged.id = len(buffer) + buffer.append(merged) + + # Keep compact buffer + if len(buffer) > 16: + _compute_pareto_counts(buffer) + buffer.sort(key=lambda c: (c.wins, c.mean), reverse=True) + buffer[:] = buffer[:16] + + # Log + best = max(buffer, key=lambda c: c.mean) + if self.logger: + self.logger.log("GEPA(base) best mean", best.mean, it, color="green") + + # Optional save + if save_path and it % 10 == 0: + _apply_params(self.optimizer, best.params) + self.save_agent(save_path, it) + + # Load best into agent + best = max(buffer, key=lambda c: c.mean) + _apply_params(self.optimizer, best.params) + return {"best_mean": best.mean}, float(best.mean) + diff --git a/tests/llm_optimizers_tests/test_gepa_benchmark.py b/tests/llm_optimizers_tests/test_gepa_benchmark.py new file mode 100644 index 00000000..fdfe5d2e --- /dev/null +++ b/tests/llm_optimizers_tests/test_gepa_benchmark.py @@ -0,0 +1,94 @@ +import os +import pytest +import numpy as np + +from opto import trace +from opto.optimizers.optoprime_v2 import OptoPrimeV2 +from opto.trainer.algorithms.gepa_algorithms import GEPAAlgorithmBase, GEPAUCBSearch, GEPABeamPareto +from opto.trainer.algorithms.basic_algorithms import BasicSearchAlgorithm +from opto.trainer.guide import LLMJudge +from opto.utils.llm import LLM + + +RUN_BENCH = "1" + + +def _datasets_or_skip(): + try: + import datasets # noqa: F401 + except Exception: + pytest.skip("datasets library not available; skipping GEPA benchmark test.") + + +def _llm_env_or_skip(): + have_key = any(os.getenv(k) for k in ["OPENAI_API_KEY", "AZURE_OPENAI_API_KEY", "ANTHROPIC_API_KEY", "OAI_CONFIG_LIST"]) + if not have_key: + pytest.skip("No LLM credentials found in environment; skipping GEPA benchmark test.") + + +@trace.model +class Learner: + """Agent that calls an LLM. The only trainable variable is 'system_prompt'.""" + + def __init__(self, system_prompt: str = "You're a helpful agent", user_prompt_template: str = "Query: {message}", llm: LLM = None): + self.system_prompt = trace.node(system_prompt, trainable=True) + self.user_prompt_template = trace.node(user_prompt_template) + self.llm = llm or LLM() # default profile + + @trace.bundle() + def model(self, system_prompt: str, user_prompt_template: str, message: str) -> str: + if "{message}" not in user_prompt_template: + raise ValueError("user_prompt_template must contain '{message}'") + resp = self.llm( + messages=[ + {"role": "system", "content": system_prompt}, + {"role": "user", "content": user_prompt_template.format(message=message)}, + ] + ) + return resp.choices[0].message.content + + def forward(self, message): + return self.model(self.system_prompt, self.user_prompt_template, message) + + +@pytest.mark.skipif(not RUN_BENCH, reason="Set RUN_GEPA_BENCH=1 to run this optional benchmark test.") +def test_gepa_benchmark_gsm8k_real_llm(): + _datasets_or_skip() + _llm_env_or_skip() + + import datasets + + # Load a tiny subset of GSM8k + ds = datasets.load_dataset("openai/gsm8k", "main") + train = ds["train"][:6] + train_dataset = {"inputs": train["question"], "infos": train["answer"]} + + # Teacher/judge with a low-cost profile + guide = LLMJudge(llm=LLM(profile="cheap")) + + # Agent and optimizer (low-cost profile) + agent = Learner(llm=LLM(profile="cheap")) + optimizer = OptoPrimeV2(agent.parameters(), llm=LLM(profile="cheap")) + + algos = [ + ("GEPA-Base", GEPAAlgorithmBase(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_iters=2, train_batch_size=1, merge_every=2)), + ("GEPA-UCB", GEPAUCBSearch(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_search_iterations=2, train_batch_size=1, merge_every=2)), + ("GEPA-Beam", GEPABeamPareto(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_search_iterations=2, train_batch_size=1, merge_every=2)), + ("BasicSearch", BasicSearchAlgorithm(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_epochs=1, batch_size=1, num_proposals=2)), + ] + + results = {} + for name, algo, kwargs in algos: + if name == "BasicSearch": + # Conform to BasicSearch's interface + algo.train(guide=guide, train_dataset=train_dataset, validate_dataset=train_dataset, test_dataset=train_dataset, eval_frequency=1, num_threads=2, verbose=False, **kwargs) + results[name] = 0.0 # placeholder; evaluation is heavy and non-deterministic + else: + _, best = algo.train(guide=guide, train_dataset=train_dataset, validate_dataset=train_dataset, pareto_subset_size=4, num_threads=2, **kwargs) + results[name] = float(best) + + # Sanity check that we produced some floats for each algorithm + assert set(results.keys()) == {"GEPA-Base", "GEPA-UCB", "GEPA-Beam", "BasicSearch"} + for v in results.values(): + assert isinstance(v, float) + From 6e518a12fdc2de12f4267bd60fec648498683194 Mon Sep 17 00:00:00 2001 From: doxav Date: Mon, 22 Sep 2025 08:17:11 +0200 Subject: [PATCH 2/6] added unit test --- tests/unit_tests/test_gepa_algorithms.py | 214 +++++++++++++++++++++++ 1 file changed, 214 insertions(+) create mode 100644 tests/unit_tests/test_gepa_algorithms.py diff --git a/tests/unit_tests/test_gepa_algorithms.py b/tests/unit_tests/test_gepa_algorithms.py new file mode 100644 index 00000000..a4c42f26 --- /dev/null +++ b/tests/unit_tests/test_gepa_algorithms.py @@ -0,0 +1,214 @@ +import math +import os +import random +import re +from typing import Any, Dict, List, Tuple + +import numpy as np +import pytest + +# Provide a light stub for optional graphviz dependency to allow imports without system graphviz +import sys, types +if "graphviz" not in sys.modules: + sys.modules["graphviz"] = types.SimpleNamespace(Digraph=object) + +from opto.trace.modules import model as trace_model +from opto.trace.nodes import node as trace_node +from opto.optimizers.optoprime_v2 import OptoPrimeV2 +import pytest +from opto.trainer.algorithms.gepa_algorithms import ( + GEPAAlgorithmBase, + GEPAUCBSearch, + GEPABeamPareto, + _compute_pareto_counts, + _pareto_sample, + _uniform_merge_params, + ) +from opto.trainer.evaluators import evaluate +from opto.trainer.guide import Guide +from opto.utils.llm import DummyLLM + + +class ExactMatchGuide(Guide): + """Simple guide: score=1 if response == reference, else 0.""" + + def get_feedback(self, query: Any, response: Any, reference: Any, **kwargs): + score = float(response == reference) + feedback = f"Score: {score}. Response: {response}. Reference: {reference}." + return score, feedback + + +@trace_model +class AddAgent: + """Toy agent: returns x + param.""" + + def __init__(self, param: int = 0): + self.param = trace_node(int(param), trainable=True) + + def forward(self, x: int) -> int: + return x + self.param + + +def make_dummy_llm(suggest_value: int) -> DummyLLM: + """Dummy LLM that parses the variable name from the prompt and suggests a fixed value. + + Matches the default XML-like output format expected by OptoPrimeV2. + """ + + def _llm_callable(messages, **kwargs): + # Extract the variable name from the #Variables section in the prompt + problem = messages[1]["content"] if isinstance(messages, (list, tuple)) and len(messages) > 1 else "" + name_match = re.findall(r"", problem) + var_name = name_match[0] if name_match else "param" + return ( + f""" + Dummy reasoning based on the input messages. + + {var_name} + {suggest_value} + + """ + ) + + return DummyLLM(_llm_callable) + + +def make_dataset(target_add: int, n: int = 8) -> Dict[str, List[int]]: + xs = list(range(n)) + infos = [x + target_add for x in xs] + return {"inputs": xs, "infos": infos} + + +def build_optimizer(agent: AddAgent, suggest_value: int) -> OptoPrimeV2: + return OptoPrimeV2(agent.parameters(), llm=make_dummy_llm(suggest_value)) + + +def test_pareto_counting_and_sampling(): + # Construct mock candidates with per-instance eval vectors where each wins on one dimension + from types import SimpleNamespace + + class Cand(SimpleNamespace): + pass + + A = Cand(eval_vector=[1.0, 0.1], wins=0, mean=0.55) + B = Cand(eval_vector=[0.2, 1.1], wins=0, mean=0.65) + cands = [A, B] + + _compute_pareto_counts(cands) + assert A.wins == 1 and B.wins == 1 + + rng = random.Random(0) + # With equal wins, both should be sampled with similar probability + picks = [ + _pareto_sample([A, B], temperature=1.0, rng=rng) for _ in range(100) + ] + a_count = sum(p is A for p in picks) + b_count = sum(p is B for p in picks) + assert abs(a_count - b_count) < 40 # rough balance + + +def test_uniform_merge_params_uses_both_parents(): + # Use two ParameterNodes to exercise merging across keys + @trace_model + class TwoParam: + def __init__(self): + self.a = trace_node(1, trainable=True) + self.b = trace_node(2, trainable=True) + + def forward(self, x): + return self.a + self.b + x + + m = TwoParam() + a_params = {p: (10 if p.py_name.endswith("a") else 20) for p in m.parameters()} + b_params = {p: (100 if p.py_name.endswith("a") else 200) for p in m.parameters()} + + rng = random.Random(123) + merged = _uniform_merge_params(a_params, b_params, rng) + # For each key, merged value should be chosen from either a_params or b_params + for k, v in merged.items(): + assert v in (a_params[k], b_params[k]) + + +@pytest.mark.parametrize( + "algo_cls,train_kwargs", + [ + (GEPAAlgorithmBase, {"num_iters": 8, "train_batch_size": 2, "merge_every": 2}), + (GEPAUCBSearch, {"num_search_iterations": 8, "train_batch_size": 2, "merge_every": 2}), + (GEPABeamPareto, {"num_search_iterations": 8, "train_batch_size": 2, "merge_every": 2}), + ], +) +def test_gepa_variants_converge_on_dummyllm(algo_cls, train_kwargs): + target_add = 5 + ds = make_dataset(target_add, n=6) + agent = AddAgent(param=0) + optimizer = build_optimizer(agent, suggest_value=target_add) + + algo = algo_cls(agent=agent, optimizer=optimizer, logger=None, num_threads=1) + + # Prepare kwargs and include 'verbose' only if supported + import inspect + call_kwargs = dict(guide=ExactMatchGuide(), train_dataset=ds, pareto_subset_size=4, num_threads=1) + sig = inspect.signature(algo.train) + if 'validation_dataset' in sig.parameters: + call_kwargs['validation_dataset'] = ds + else: + call_kwargs['validate_dataset'] = ds + call_kwargs.update(train_kwargs) + if 'verbose' in sig.parameters: + call_kwargs['verbose'] = False + + metrics, best = algo.train(**call_kwargs) + + # Best mean on pareto subset should be perfect + assert isinstance(best, float) + assert best == pytest.approx(1.0, rel=0, abs=1e-6) + # Agent parameter should be updated to target_add + assert agent.param.data == target_add + + +def test_compare_gepa_vs_basicsearch_on_dummyllm(): + from opto.trainer.algorithms.basic_algorithms import BasicSearchAlgorithm + + target_add = 7 + ds = make_dataset(target_add, n=6) + agent_gepa = AddAgent(param=0) + agent_basic = AddAgent(param=0) + + opt_gepa = build_optimizer(agent_gepa, suggest_value=target_add) + opt_basic = build_optimizer(agent_basic, suggest_value=target_add) + + # GEPA + gepa = GEPAAlgorithmBase(agent_gepa, optimizer=opt_gepa, logger=None, num_threads=1) + _, best_gepa = gepa.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + pareto_subset_size=4, + num_iters=8, + train_batch_size=2, + merge_every=2, + num_threads=1, + ) + + # BasicSearch baseline + basic = BasicSearchAlgorithm(agent_basic, optimizer=opt_basic, logger=None, num_threads=1) + basic.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + num_proposals=1, + num_epochs=1, + batch_size=1, + test_dataset=ds, + eval_frequency=1, + num_threads=1, + verbose=False, + ) + + # Evaluate both on full dataset + score_gepa = np.mean(evaluate(agent_gepa, ExactMatchGuide(), ds["inputs"], ds["infos"], num_threads=2)) + score_basic = np.mean(evaluate(agent_basic, ExactMatchGuide(), ds["inputs"], ds["infos"], num_threads=2)) + + assert best_gepa == pytest.approx(1.0, rel=0, abs=1e-6) + assert score_gepa == pytest.approx(1.0, rel=0, abs=1e-6) + assert score_basic == pytest.approx(1.0, rel=0, abs=1e-6) From d8b7269632ca93dfa7be24beae560cdaf96ba45a Mon Sep 17 00:00:00 2001 From: doxav Date: Mon, 22 Sep 2025 08:52:00 +0200 Subject: [PATCH 3/6] IMPROVED but to check: total iterations seems much higher --- opto/trainer/algorithms/gepa_algorithms.py | 246 ++++++++++++++--- .../test_gepa_benchmark.py | 64 ++++- tests/unit_tests/test_gepa_algorithms.py | 257 ++++++++++++++++++ 3 files changed, 521 insertions(+), 46 deletions(-) diff --git a/opto/trainer/algorithms/gepa_algorithms.py b/opto/trainer/algorithms/gepa_algorithms.py index 588cdbad..2dac2a7c 100644 --- a/opto/trainer/algorithms/gepa_algorithms.py +++ b/opto/trainer/algorithms/gepa_algorithms.py @@ -26,6 +26,11 @@ standard_optimization_step, ) from opto.trainer.utils import async_run +# Prefer thread-safe batched runner (deep-copies per task). Fallback handled at callsite. +try: + from opto.trainer.utils import batch_run # type: ignore +except Exception: # pragma: no cover + batch_run = None from opto.optimizers.utils import print_color @@ -134,7 +139,7 @@ def _maybe_merge(buffer: List[Candidate], merged_params = _uniform_merge_params(a.params, b.params, rng) # Evaluate merged on Pareto subset - original_params = {p: copy.deepcopy(p.data) for p in agent.parameters()} + original_params = _snapshot_params_fast(list(agent.parameters())) try: # load params to agent from opto.optimizers.optimizer import Optimizer # type: ignore @@ -163,6 +168,98 @@ def _maybe_merge(buffer: List[Candidate], return None +def _maybe_merge_ancestor_aware( + buffer: List[Candidate], + *, + id2cand: Dict[int, Candidate], + module_groups: List[List[ParameterNode]], + agent, + guide, + optimizer, + train_dataset: Dict[str, List[Any]], + train_batch_size: int, + pareto_inputs: List[Any], + pareto_infos: List[Any], + num_threads: Optional[int], + rng: random.Random, + tried_pairs: set, + budget_tracker: Optional[Dict[str, int]] = None, + budget_B: Optional[int] = None, + max_tries: int = 8 +) -> Optional[Tuple[Candidate, int]]: + """ + Ancestor-aware merge with budget tracking. Returns (merged_candidate, rollouts_used). + """ + if len(buffer) < 2: + return None + + rollouts_used = 0 + + # Sample training minibatch + tx = rng.choices(train_dataset["inputs"], k=min(train_batch_size, len(train_dataset["inputs"]))) + ti = rng.choices(train_dataset["infos"], k=len(tx)) + + # Prefer winners for parent selection + _compute_pareto_counts(buffer) + pool = sorted(buffer, key=lambda c: (c.wins, c.mean), reverse=True) + + for _ in range(max_tries): + i, j = rng.sample(range(len(pool)), 2) + ci, cj = pool[i], pool[j] + if ci.id == cj.id: + continue + if ci.id in cj.ancestors or cj.id in ci.ancestors: + continue # avoid direct ancestry + key = tuple(sorted((ci.id, cj.id))) + if key in tried_pairs: + continue + tried_pairs.add(key) + + merged_params = _uniform_merge_params(ci.params, cj.params, rng) + + # Quick minibatch acceptability check + def _batch_mean_for(param_dict): + original = _snapshot_params_fast(list(optimizer.parameters)) + try: + _apply_params(optimizer, param_dict) + vec = evaluate(agent, guide, tx, ti, min_score=None, num_threads=num_threads, + description="MERGE(mini-batch accept)") + finally: + _apply_params(optimizer, original) + return float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + + rollouts_used += len(tx) + merged_batch_mean = _batch_mean_for(merged_params) + parent_means = [_batch_mean_for(ci.params), _batch_mean_for(cj.params)] + rollouts_used += 2 * len(tx) + + if merged_batch_mean <= max(parent_means): + continue # Not promising enough + + # Full Pareto evaluation + original = _snapshot_params_fast(list(optimizer.parameters)) + try: + _apply_params(optimizer, merged_params) + vec = evaluate(agent, guide, pareto_inputs, pareto_infos, min_score=None, + num_threads=num_threads, description="GEPA+Merge: ancestor-aware Pareto eval") + finally: + _apply_params(optimizer, original) + mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") + # Account Pareto evaluation cost in the global budget and local counter. + if budget_B is not None and budget_tracker is not None: + budget_tracker["used"] += len(pareto_inputs) + rollouts_used += len(pareto_inputs) + + merged = Candidate(params=merged_params, + eval_vector=vec, mean=mean, + id=-1, parent_ids=(ci.id, cj.id), + ancestors=set(ci.ancestors) | set(cj.ancestors) | {ci.id, cj.id}, + created_iter=0) + return merged, rollouts_used + + return None + + def _ensure_optimizer(agent, optimizer): if optimizer is not None: return optimizer @@ -175,14 +272,16 @@ def _train_step_generate_child(agent, guide, optimizer, train_xs, train_infos, * Single-parent, incremental evolution "mutation": run forward on a minibatch to get batched feedback, then optimizer.step(bypassing=True) to obtain a new candidate param dict (without applying). """ - use_async = num_threads is not None and num_threads > 1 - if use_async: + use_parallel = (num_threads is not None and num_threads > 1) + if use_parallel: + # Use async_run but ensure thread safety through parameter handling + # Since we're working with parameters through optimizer, this should be thread-safe outputs = async_run([lambda a,x,g,info: standard_optimization_step(a, x, g, info)] * len(train_xs), args_list=[(agent, x, guide, info) for x, info in zip(train_xs, train_infos)], max_workers=num_threads, description="GEPA forward (mutate parent)") - # outputs: List[(target, score, feedback)] else: + # Safe sequential fallback. outputs = [standard_optimization_step(agent, x, guide, info) for x, info in zip(train_xs, train_infos)] scores, targets, feedbacks = [], [], [] @@ -212,6 +311,44 @@ def _apply_params(optimizer, param_dict: Dict[ParameterNode, Any]): optimizer.update(param_dict) +def _snapshot_params_fast(parameters: List[ParameterNode]) -> Dict[ParameterNode, Any]: + """ + Snapshot ParameterNode->value with minimal copying: + - immutables (str/int/float/bool/tuple/bytes/None): no copy + - numpy arrays: .copy() + - everything else: deepcopy (safe fallback) + """ + snap: Dict[ParameterNode, Any] = {} + immutables = (str, int, float, bool, tuple, frozenset, bytes, type(None)) + for p in parameters: + v = getattr(p, "data", None) + if isinstance(v, immutables): + snap[p] = v + elif isinstance(v, np.ndarray): + snap[p] = v.copy() + else: + snap[p] = copy.deepcopy(v) + return snap + + +def _fingerprint_params(params_dict: Dict[ParameterNode, Any]) -> Tuple: + """ + Hashable fingerprint of a ParameterNode->value dict for optional caching. + Uses (param-id, repr(value)) with special handling for numpy arrays. + """ + items: List[Tuple] = [] + for p, v in params_dict.items(): + pid = getattr(p, "uid", None) or getattr(p, "name", None) or id(p) + try: + if isinstance(v, np.ndarray): + items.append(("arr", pid, v.shape, v.dtype.str, hash(v.tobytes()))) + else: + items.append(("val", pid, repr(v))) + except Exception: + items.append(("val", pid, repr(v))) + return tuple(sorted(items)) + + # ======================= Variant 1: GEPA + Merge (UCB subclass) ======================= # class GEPAUCBSearch(UCBSearchAlgorithm): @@ -232,7 +369,10 @@ def __init__(self, ucb_exploration_factor: float = 0.8, rng_seed: int = 7, logger=None, - num_threads: Optional[int] = None): + num_threads: Optional[int] = None, + module_groups: Optional[Dict[str, List[ParameterNode]] | List[List[ParameterNode]]] = None, + selectmodule_policy: str = "round_robin", + enable_pareto_cache: bool = False): optimizer = _ensure_optimizer(agent, optimizer) super().__init__(agent, optimizer, max_buffer_size=max_buffer_size, @@ -240,22 +380,37 @@ def __init__(self, logger=logger, num_threads=num_threads) self.rng = random.Random(rng_seed) + np.random.seed(rng_seed) # ensure numpy reproducibility for np.random.choice self._pareto_inputs: List[Any] = [] self._pareto_infos: List[Any] = [] self._id_counter = 0 + self.enable_pareto_cache = enable_pareto_cache + self._pareto_cache: Dict[Tuple, Tuple[List[float], float]] = {} + # >>> NEW selector (commented out as ModuleSelector may not exist) + # self.module_selector = ModuleSelector(self.optimizer.parameters, + # module_groups=module_groups, + # policy=selectmodule_policy) def _next_id(self) -> int: self._id_counter += 1 return self._id_counter def _evaluate_on_pareto(self, params_dict: Dict[ParameterNode, Any], guide, *, num_threads) -> Tuple[List[float], float]: - original_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + cache_key = _fingerprint_params(params_dict) if self.enable_pareto_cache else None + if cache_key is not None: + cached = self._pareto_cache.get(cache_key) + if cached is not None: + return cached + original_params = _snapshot_params_fast(list(self.optimizer.parameters)) try: _apply_params(self.optimizer, params_dict) vec = _eval_on_subset(self.agent, guide, self._pareto_inputs, self._pareto_infos, num_threads=num_threads, desc="GEPA: evaluate on Pareto subset") mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - return vec, mean + result = (vec, mean) + if cache_key is not None: + self._pareto_cache[cache_key] = result + return result finally: _apply_params(self.optimizer, original_params) @@ -293,11 +448,14 @@ def train(self, buffer: List[Candidate] = [] tried_merges: set = set() + id2cand: Dict[int, Candidate] = {} # Seed with current params - base_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + base_params = _snapshot_params_fast(list(self.optimizer.parameters)) v0, m0 = self._evaluate_on_pareto(base_params, guide, num_threads=num_threads) - buffer.append(Candidate(params=base_params, eval_vector=v0, mean=m0, id=self._next_id(), ancestors=set())) + seed = Candidate(params=base_params, eval_vector=v0, mean=m0, id=self._next_id(), ancestors=set(), created_iter=0) + buffer.append(seed) + id2cand[seed.id] = seed print_color(f"[GEPA] Seed candidate mean={m0:.4f}", "cyan") metrics = {"best_means": [], "new_child_means": [], "merge_accepts": 0, "total_merges": 0} @@ -384,16 +542,17 @@ class GEPABeamPareto(BeamsearchAlgorithm): - replace deep beam expansion with GEPA’s single-parent incremental evolution """ - def __init__(self, - agent, - optimizer=None, - *, - rng_seed: int = 11, - logger=None, - num_threads: Optional[int] = None): + def __init__(self, agent, optimizer=None, *, rng_seed: int = 11, logger=None, + num_threads: Optional[int] = None, + module_groups: Optional[Dict[str, List[ParameterNode]] | List[List[ParameterNode]]] = None, + selectmodule_policy: str = "round_robin"): optimizer = _ensure_optimizer(agent, optimizer) super().__init__(agent, optimizer, num_threads=num_threads, logger=logger) self.rng = random.Random(rng_seed) + np.random.seed(rng_seed) + # self.module_selector = ModuleSelector(self.optimizer.parameters, + # module_groups=module_groups, + # policy=selectmodule_policy) # We keep a Pareto select helper that returns (selected_params, wins, scores) def select(self, @@ -409,7 +568,7 @@ def select(self, """ # Evaluate each candidate to a vector on the mini validation cand_objs: List[Candidate] = [] - current_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + current_params = _snapshot_params_fast(list(self.optimizer.parameters)) try: for idx, params in enumerate(candidates): _apply_params(self.optimizer, params) @@ -436,20 +595,17 @@ def select(self, return sel_params # Replace beam "train" with GEPA-style incremental loop (keeps BeamsearchAlgorithm API) - def train(self, - guide, - train_dataset, - *, - validate_dataset=None, - pareto_subset_size: int = 24, - num_search_iterations: int = 120, - train_batch_size: int = 2, - merge_every: int = 6, - log_frequency: Optional[int] = None, + def train(self, guide, train_dataset, *, + validate_dataset=None, pareto_subset_size: int = 24, + num_search_iterations: int = 120, train_batch_size: int = 2, + merge_every: int = 6, log_frequency: Optional[int] = None, save_frequency: Optional[int] = None, save_path: str = "checkpoints/gepa_beam_agent.pkl", - verbose: bool = False, - num_threads: Optional[int] = None): + verbose: bool = False, num_threads: Optional[int] = None, + module_groups: Optional[Dict[str, List[ParameterNode]] | List[List[ParameterNode]]] = None, + selectmodule_policy: str = "round_robin", + budget_B: Optional[int] = None, + accept_epsilon: float = 0.0): num_threads = num_threads or self.num_threads log_frequency = log_frequency or 5 validate_ds = validate_dataset or train_dataset @@ -463,16 +619,15 @@ def train(self, # Seed buffer buffer: List[Candidate] = [] - base_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} - # Evaluate seed - current_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + base_params = _snapshot_params_fast(list(self.optimizer.parameters)) + original = _snapshot_params_fast(list(self.optimizer.parameters)) try: _apply_params(self.optimizer, base_params) vec = evaluate(self.agent, guide, pareto_inputs, pareto_infos, min_score=None, num_threads=num_threads, description="GEPA(beam): seed evaluation") finally: - _apply_params(self.optimizer, current_params) + _apply_params(self.optimizer, original) m0 = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") buffer.append(Candidate(params=base_params, eval_vector=vec, mean=m0, id=0, ancestors=set())) tried_merges: set = set() @@ -496,13 +651,13 @@ def train(self, continue # Evaluate child on Pareto subset - current_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + original = _snapshot_params_fast(list(self.optimizer.parameters)) try: _apply_params(self.optimizer, update_dict) vec = evaluate(self.agent, guide, pareto_inputs, pareto_infos, min_score=None, num_threads=num_threads, description="GEPA(beam): child eval") finally: - _apply_params(self.optimizer, current_params) + _apply_params(self.optimizer, original) mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") buffer.append(Candidate(params=update_dict, eval_vector=vec, mean=mean, id=len(buffer), parent_ids=(parent.id,), ancestors=set(parent.ancestors) | {parent.id})) @@ -545,16 +700,17 @@ class GEPAAlgorithmBase(Trainer): Useful when you want the simplest control loop with your own logging/saving. """ - def __init__(self, - agent, - optimizer=None, - *, - rng_seed: int = 13, - logger=None, - num_threads: Optional[int] = None): + def __init__(self, agent, optimizer=None, *, rng_seed: int = 13, logger=None, + num_threads: Optional[int] = None, + module_groups: Optional[Dict[str, List[ParameterNode]] | List[List[ParameterNode]]] = None, + selectmodule_policy: str = "round_robin"): super().__init__(agent, num_threads=num_threads, logger=logger) self.optimizer = _ensure_optimizer(agent, optimizer) self.rng = random.Random(rng_seed) + np.random.seed(rng_seed) + # self.module_selector = ModuleSelector(self.optimizer.parameters, + # module_groups=module_groups, + # policy=selectmodule_policy) def train(self, guide, @@ -579,8 +735,8 @@ def train(self, # Seed buffer: List[Candidate] = [] - base_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} - original = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + base_params = _snapshot_params_fast(list(self.optimizer.parameters)) + original = _snapshot_params_fast(list(self.optimizer.parameters)) try: _apply_params(self.optimizer, base_params) vec = evaluate(self.agent, guide, xsP, isP, min_score=None, num_threads=num_threads, @@ -608,7 +764,7 @@ def train(self, continue # Eval child - original = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} + original = _snapshot_params_fast(list(self.optimizer.parameters)) try: _apply_params(self.optimizer, update_dict) vec = evaluate(self.agent, guide, xsP, isP, min_score=None, num_threads=num_threads, diff --git a/tests/llm_optimizers_tests/test_gepa_benchmark.py b/tests/llm_optimizers_tests/test_gepa_benchmark.py index fdfe5d2e..a9254725 100644 --- a/tests/llm_optimizers_tests/test_gepa_benchmark.py +++ b/tests/llm_optimizers_tests/test_gepa_benchmark.py @@ -66,14 +66,18 @@ def test_gepa_benchmark_gsm8k_real_llm(): # Teacher/judge with a low-cost profile guide = LLMJudge(llm=LLM(profile="cheap")) + # Set a budget constraint for algorithms that support it (e.g., GEPABeamPareto) + budget_limit = 5 + # Agent and optimizer (low-cost profile) agent = Learner(llm=LLM(profile="cheap")) optimizer = OptoPrimeV2(agent.parameters(), llm=LLM(profile="cheap")) algos = [ ("GEPA-Base", GEPAAlgorithmBase(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_iters=2, train_batch_size=1, merge_every=2)), + (f"GEPA-BeamPareto-Budget{budget_limit}", GEPABeamPareto(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_search_iterations=2, train_batch_size=1, merge_every=2, budget_B=budget_limit)), + ("GEPA-BeamPareto", GEPABeamPareto(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_search_iterations=2, train_batch_size=1, merge_every=2)), ("GEPA-UCB", GEPAUCBSearch(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_search_iterations=2, train_batch_size=1, merge_every=2)), - ("GEPA-Beam", GEPABeamPareto(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_search_iterations=2, train_batch_size=1, merge_every=2)), ("BasicSearch", BasicSearchAlgorithm(agent, optimizer=optimizer, logger=None, num_threads=2), dict(num_epochs=1, batch_size=1, num_proposals=2)), ] @@ -92,3 +96,61 @@ def test_gepa_benchmark_gsm8k_real_llm(): for v in results.values(): assert isinstance(v, float) + +# @pytest.mark.skipif(not RUN_BENCH, reason="Set RUN_GEPA_BENCH=1 to run this optional benchmark test.") +# def test_gepa_benchmark_gsm8k_low_budget(): +# """Same benchmark test but with a low budget constraint (5 evaluations).""" +# _datasets_or_skip() +# _llm_env_or_skip() + +# import datasets + +# # Load a tiny subset of GSM8k +# ds = datasets.load_dataset("openai/gsm8k", "main") +# train = ds["train"][:3] # Even smaller dataset for low budget +# train_dataset = {"inputs": train["question"], "infos": train["answer"]} + +# # Teacher/judge with a low-cost profile +# guide = LLMJudge(llm=LLM(profile="cheap")) + +# # Test each GEPA variant with budget constraint +# budget_limit = 5 +# algos = [ +# ("GEPA-Base-Budget", GEPAAlgorithmBase, dict(num_iters=1, train_batch_size=1, merge_every=2)), +# ("GEPA-UCB-Budget", GEPAUCBSearch, dict(num_search_iterations=1, train_batch_size=1, merge_every=2)), +# ("GEPA-Beam-Budget", GEPABeamPareto, dict(num_search_iterations=1, train_batch_size=1, merge_every=2, budget_B=budget_limit)), +# ] + +# results = {} +# for name, algo_cls, kwargs in algos: +# # Create fresh agent and optimizer for each test +# agent = Learner(llm=LLM(profile="cheap")) +# optimizer = OptoPrimeV2(agent.parameters(), llm=LLM(profile="cheap")) +# algo = algo_cls(agent, optimizer=optimizer, logger=None, num_threads=1) + +# # Add budget_B to kwargs if supported by the algorithm +# if name == "GEPA-Beam-Budget": +# # GEPABeamPareto supports budget_B parameter +# pass # budget_B already in kwargs + +# try: +# _, best = algo.train( +# guide=guide, +# train_dataset=train_dataset, +# validate_dataset=train_dataset, +# pareto_subset_size=2, # Small Pareto subset to save budget +# num_threads=1, +# **kwargs +# ) +# results[name] = float(best) +# except Exception as e: +# # If budget constraint causes early termination or other issues, record as 0 +# print(f"Algorithm {name} encountered error with budget constraint: {e}") +# results[name] = 0.0 + +# # Sanity check that we produced some results +# assert set(results.keys()) == {"GEPA-Base-Budget", "GEPA-UCB-Budget", "GEPA-Beam-Budget"} +# for v in results.values(): +# assert isinstance(v, float) +# assert v >= 0.0 # Should be non-negative scores + diff --git a/tests/unit_tests/test_gepa_algorithms.py b/tests/unit_tests/test_gepa_algorithms.py index a4c42f26..c435628e 100644 --- a/tests/unit_tests/test_gepa_algorithms.py +++ b/tests/unit_tests/test_gepa_algorithms.py @@ -212,3 +212,260 @@ def test_compare_gepa_vs_basicsearch_on_dummyllm(): assert best_gepa == pytest.approx(1.0, rel=0, abs=1e-6) assert score_gepa == pytest.approx(1.0, rel=0, abs=1e-6) assert score_basic == pytest.approx(1.0, rel=0, abs=1e-6) + + +def test_snapshot_params_fast(): + """Test the fast parameter snapshot utility function.""" + from opto.trainer.algorithms.gepa_algorithms import _snapshot_params_fast + + @trace_model + class MultiTypeAgent: + def __init__(self): + self.int_param = trace_node(42, trainable=True) + self.str_param = trace_node("hello", trainable=True) + self.float_param = trace_node(3.14, trainable=True) + self.list_param = trace_node([1, 2, 3], trainable=True) + self.dict_param = trace_node({"key": "value"}, trainable=True) + # Test numpy array + self.np_param = trace_node(np.array([1, 2, 3]), trainable=True) + + def forward(self, x): + return x + self.int_param + + agent = MultiTypeAgent() + params = list(agent.parameters()) + + # Test snapshot + snapshot = _snapshot_params_fast(params) + + # Check that all parameters are included + assert len(snapshot) == len(params) + + # Modify original values + agent.int_param._set(100) + agent.str_param._set("modified") + agent.np_param._set(np.array([4, 5, 6])) + + # Verify snapshot preserved original values + for p in params: + if p.py_name == "int_param": + assert snapshot[p] == 42 + elif p.py_name == "str_param": + assert snapshot[p] == "hello" + elif p.py_name == "np_param": + assert np.array_equal(snapshot[p], np.array([1, 2, 3])) + + +def test_fingerprint_params(): + """Test the parameter fingerprinting utility function.""" + from opto.trainer.algorithms.gepa_algorithms import _fingerprint_params + + @trace_model + class SimpleAgent: + def __init__(self): + self.a = trace_node(1, trainable=True) + self.b = trace_node("test", trainable=True) + + def forward(self, x): + return x + self.a + + agent = SimpleAgent() + params_dict = {p: p.data for p in agent.parameters()} + + # Test fingerprinting + fp1 = _fingerprint_params(params_dict) + fp2 = _fingerprint_params(params_dict) + + # Same parameters should produce same fingerprint + assert fp1 == fp2 + + # Different parameters should produce different fingerprint + agent.a._set(2) + params_dict2 = {p: p.data for p in agent.parameters()} + fp3 = _fingerprint_params(params_dict2) + assert fp1 != fp3 + + +def test_numpy_seeding_reproducibility(): + """Test that numpy seeding ensures reproducible behavior.""" + target_add = 3 + ds = make_dataset(target_add, n=4) + + # Test with same seed + results = [] + for seed in [123, 123]: # Same seed twice + agent = AddAgent(param=0) + optimizer = build_optimizer(agent, suggest_value=target_add) + algo = GEPAAlgorithmBase(agent=agent, optimizer=optimizer, logger=None, num_threads=1, rng_seed=seed) + + metrics, best = algo.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + pareto_subset_size=3, + num_iters=2, + train_batch_size=1, + merge_every=2, + num_threads=1, + ) + results.append((metrics, best, agent.param.data)) + + # Results should be identical with same seed + assert results[0][1] == results[1][1] # Same best score + assert results[0][2] == results[1][2] # Same final parameter + + # Test with different seed + agent_diff = AddAgent(param=0) + optimizer_diff = build_optimizer(agent_diff, suggest_value=target_add) + algo_diff = GEPAAlgorithmBase(agent=agent_diff, optimizer=optimizer_diff, logger=None, num_threads=1, rng_seed=456) + + metrics_diff, best_diff = algo_diff.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + pareto_subset_size=3, + num_iters=2, + train_batch_size=1, + merge_every=2, + num_threads=1, + ) + + # Both should converge but the process might differ + # (though with DummyLLM behavior is very predictable) + assert best_diff == pytest.approx(1.0, rel=0, abs=1e-6) + + +def test_gepa_ucb_pareto_cache(): + """Test Pareto cache functionality in GEPAUCBSearch.""" + target_add = 4 + ds = make_dataset(target_add, n=3) + agent = AddAgent(param=0) + optimizer = build_optimizer(agent, suggest_value=target_add) + + # Test with cache enabled + algo = GEPAUCBSearch(agent=agent, optimizer=optimizer, logger=None, num_threads=1, enable_pareto_cache=True) + + metrics, best = algo.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + pareto_subset_size=2, + num_search_iterations=2, + train_batch_size=1, + merge_every=2, + num_threads=1, + ) + + # Should converge to perfect solution + assert best == pytest.approx(1.0, rel=0, abs=1e-6) + assert agent.param.data == target_add + + # Test that cache was used (should have some entries) + # Note: exact cache size depends on algorithm behavior, but should be non-empty if enabled + if hasattr(algo, '_pareto_cache'): + assert isinstance(algo._pareto_cache, dict) + + +def test_budget_tracking_functionality(): + """Test budget tracking in GEPA algorithms.""" + target_add = 2 + ds = make_dataset(target_add, n=4) + agent = AddAgent(param=0) + optimizer = build_optimizer(agent, suggest_value=target_add) + + # Test GEPABeamPareto with budget + algo = GEPABeamPareto(agent=agent, optimizer=optimizer, logger=None, num_threads=1) + + metrics, best = algo.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + pareto_subset_size=3, + num_search_iterations=2, + train_batch_size=1, + merge_every=2, + budget_B=10, # Low budget to test tracking + num_threads=1, + ) + + # Should still achieve good results even with budget constraint + assert isinstance(best, float) + assert best >= 0.0 # Should be non-negative score + + +def test_thread_safety_with_sequential_fallback(): + """Test that algorithms work correctly with sequential fallback when batch_run unavailable.""" + target_add = 1 + ds = make_dataset(target_add, n=2) + agent = AddAgent(param=0) + optimizer = build_optimizer(agent, suggest_value=target_add) + + # Test with num_threads=1 (should use sequential) + algo = GEPAAlgorithmBase(agent=agent, optimizer=optimizer, logger=None, num_threads=1) + metrics, best = algo.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + pareto_subset_size=2, + num_iters=2, + train_batch_size=1, + merge_every=2, + num_threads=1, + ) + + assert best == pytest.approx(1.0, rel=0, abs=1e-6) + assert agent.param.data == target_add + + # Test with num_threads=2 (may use parallel or fallback to sequential) + agent2 = AddAgent(param=0) + optimizer2 = build_optimizer(agent2, suggest_value=target_add) + algo2 = GEPAAlgorithmBase(agent=agent2, optimizer=optimizer2, logger=None, num_threads=2) + + metrics2, best2 = algo2.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + pareto_subset_size=2, + num_iters=2, + train_batch_size=1, + merge_every=2, + num_threads=2, + ) + + assert best2 == pytest.approx(1.0, rel=0, abs=1e-6) + assert agent2.param.data == target_add + + +def test_gepa_ucb_selectmodule_policy(): + """Test different module selection policies in GEPAUCBSearch.""" + target_add = 6 + ds = make_dataset(target_add, n=3) + + # Test different selection policies + policies = ["round_robin"] # Could test more if other policies are available + + for policy in policies: + agent = AddAgent(param=0) + optimizer = build_optimizer(agent, suggest_value=target_add) + + algo = GEPAUCBSearch( + agent=agent, + optimizer=optimizer, + logger=None, + num_threads=1, + selectmodule_policy=policy + ) + + metrics, best = algo.train( + guide=ExactMatchGuide(), + train_dataset=ds, + validate_dataset=ds, + pareto_subset_size=2, + num_search_iterations=2, + train_batch_size=1, + merge_every=2, + num_threads=1, + ) + + assert best == pytest.approx(1.0, rel=0, abs=1e-6) + assert agent.param.data == target_add From 0b306faa0b841e14657d6ec04a4a15a2b7f60133 Mon Sep 17 00:00:00 2001 From: doxav Date: Mon, 22 Sep 2025 21:15:06 +0200 Subject: [PATCH 4/6] added budget and stabilized parallel training --- opto/trainer/algorithms/gepa_algorithms.py | 77 +++++++++++++------ .../test_gepa_benchmark.py | 58 -------------- 2 files changed, 55 insertions(+), 80 deletions(-) diff --git a/opto/trainer/algorithms/gepa_algorithms.py b/opto/trainer/algorithms/gepa_algorithms.py index 2dac2a7c..c0283e38 100644 --- a/opto/trainer/algorithms/gepa_algorithms.py +++ b/opto/trainer/algorithms/gepa_algorithms.py @@ -10,6 +10,8 @@ import copy import math import random +import functools +import types from dataclasses import dataclass, field from typing import Any, Dict, List, Optional, Tuple @@ -195,9 +197,11 @@ def _maybe_merge_ancestor_aware( rollouts_used = 0 - # Sample training minibatch - tx = rng.choices(train_dataset["inputs"], k=min(train_batch_size, len(train_dataset["inputs"]))) - ti = rng.choices(train_dataset["infos"], k=len(tx)) + # Sample training minibatch (no replacement → lower variance) + k = min(train_batch_size, len(train_dataset["inputs"])) + idxs = np.random.choice(len(train_dataset["inputs"]), k, replace=False) + tx = [train_dataset["inputs"][i] for i in idxs] + ti = [train_dataset["infos"][i] for i in idxs] # Prefer winners for parent selection _compute_pareto_counts(buffer) @@ -228,10 +232,16 @@ def _batch_mean_for(param_dict): _apply_params(optimizer, original) return float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - rollouts_used += len(tx) + rollouts_used += k merged_batch_mean = _batch_mean_for(merged_params) parent_means = [_batch_mean_for(ci.params), _batch_mean_for(cj.params)] - rollouts_used += 2 * len(tx) + rollouts_used += 2 * k + + # Early budget guard (3*k minibatch evals) before Pareto eval + if budget_B is not None and budget_tracker is not None: + if budget_tracker["used"] + 3 * k + len(pareto_inputs) > budget_B: + return None + budget_tracker["used"] += 3 * k if merged_batch_mean <= max(parent_means): continue # Not promising enough @@ -272,14 +282,47 @@ def _train_step_generate_child(agent, guide, optimizer, train_xs, train_infos, * Single-parent, incremental evolution "mutation": run forward on a minibatch to get batched feedback, then optimizer.step(bypassing=True) to obtain a new candidate param dict (without applying). """ - use_parallel = (num_threads is not None and num_threads > 1) + use_parallel = (num_threads is not None and num_threads > 1 and batch_run is not None) if use_parallel: - # Use async_run but ensure thread safety through parameter handling - # Since we're working with parameters through optimizer, this should be thread-safe - outputs = async_run([lambda a,x,g,info: standard_optimization_step(a, x, g, info)] * len(train_xs), - args_list=[(agent, x, guide, info) for x, info in zip(train_xs, train_infos)], - max_workers=num_threads, - description="GEPA forward (mutate parent)") + # Pre-bind args → pass callables only. Robust to different batch_run signatures. + callables = [ + functools.partial(standard_optimization_step, agent, x, guide, info) + for x, info in zip(train_xs, train_infos) + ] + try: + outputs = batch_run( + callables, + max_workers=num_threads, + description="GEPA forward (mutate parent)", + ) + except TypeError: + # Fallback: older/other signature (e.g., batch_run(callables, max_workers)) + try: + outputs = batch_run(callables, num_threads) + except Exception: + outputs = None + # Normalize outputs to a list of results. batch_run in different versions may: + # - return the list of results, + # - return a callable that returns the results, + # - return a generator/iterator, + # - or return None. + try: + if callable(outputs): + outputs = outputs() + elif isinstance(outputs, types.GeneratorType): + outputs = list(outputs) + elif outputs is None: + # fallback to sequential evaluation + outputs = [fn() for fn in callables] + elif not isinstance(outputs, (list, tuple)): + # Some other iterable (e.g. map object) + try: + outputs = list(outputs) + except Exception: + outputs = [fn() for fn in callables] + except Exception: + # Any error while normalizing → fallback to sequential + outputs = [fn() for fn in callables] else: # Safe sequential fallback. outputs = [standard_optimization_step(agent, x, guide, info) for x, info in zip(train_xs, train_infos)] @@ -386,10 +429,6 @@ def __init__(self, self._id_counter = 0 self.enable_pareto_cache = enable_pareto_cache self._pareto_cache: Dict[Tuple, Tuple[List[float], float]] = {} - # >>> NEW selector (commented out as ModuleSelector may not exist) - # self.module_selector = ModuleSelector(self.optimizer.parameters, - # module_groups=module_groups, - # policy=selectmodule_policy) def _next_id(self) -> int: self._id_counter += 1 @@ -550,9 +589,6 @@ def __init__(self, agent, optimizer=None, *, rng_seed: int = 11, logger=None, super().__init__(agent, optimizer, num_threads=num_threads, logger=logger) self.rng = random.Random(rng_seed) np.random.seed(rng_seed) - # self.module_selector = ModuleSelector(self.optimizer.parameters, - # module_groups=module_groups, - # policy=selectmodule_policy) # We keep a Pareto select helper that returns (selected_params, wins, scores) def select(self, @@ -708,9 +744,6 @@ def __init__(self, agent, optimizer=None, *, rng_seed: int = 13, logger=None, self.optimizer = _ensure_optimizer(agent, optimizer) self.rng = random.Random(rng_seed) np.random.seed(rng_seed) - # self.module_selector = ModuleSelector(self.optimizer.parameters, - # module_groups=module_groups, - # policy=selectmodule_policy) def train(self, guide, diff --git a/tests/llm_optimizers_tests/test_gepa_benchmark.py b/tests/llm_optimizers_tests/test_gepa_benchmark.py index a9254725..31efa305 100644 --- a/tests/llm_optimizers_tests/test_gepa_benchmark.py +++ b/tests/llm_optimizers_tests/test_gepa_benchmark.py @@ -96,61 +96,3 @@ def test_gepa_benchmark_gsm8k_real_llm(): for v in results.values(): assert isinstance(v, float) - -# @pytest.mark.skipif(not RUN_BENCH, reason="Set RUN_GEPA_BENCH=1 to run this optional benchmark test.") -# def test_gepa_benchmark_gsm8k_low_budget(): -# """Same benchmark test but with a low budget constraint (5 evaluations).""" -# _datasets_or_skip() -# _llm_env_or_skip() - -# import datasets - -# # Load a tiny subset of GSM8k -# ds = datasets.load_dataset("openai/gsm8k", "main") -# train = ds["train"][:3] # Even smaller dataset for low budget -# train_dataset = {"inputs": train["question"], "infos": train["answer"]} - -# # Teacher/judge with a low-cost profile -# guide = LLMJudge(llm=LLM(profile="cheap")) - -# # Test each GEPA variant with budget constraint -# budget_limit = 5 -# algos = [ -# ("GEPA-Base-Budget", GEPAAlgorithmBase, dict(num_iters=1, train_batch_size=1, merge_every=2)), -# ("GEPA-UCB-Budget", GEPAUCBSearch, dict(num_search_iterations=1, train_batch_size=1, merge_every=2)), -# ("GEPA-Beam-Budget", GEPABeamPareto, dict(num_search_iterations=1, train_batch_size=1, merge_every=2, budget_B=budget_limit)), -# ] - -# results = {} -# for name, algo_cls, kwargs in algos: -# # Create fresh agent and optimizer for each test -# agent = Learner(llm=LLM(profile="cheap")) -# optimizer = OptoPrimeV2(agent.parameters(), llm=LLM(profile="cheap")) -# algo = algo_cls(agent, optimizer=optimizer, logger=None, num_threads=1) - -# # Add budget_B to kwargs if supported by the algorithm -# if name == "GEPA-Beam-Budget": -# # GEPABeamPareto supports budget_B parameter -# pass # budget_B already in kwargs - -# try: -# _, best = algo.train( -# guide=guide, -# train_dataset=train_dataset, -# validate_dataset=train_dataset, -# pareto_subset_size=2, # Small Pareto subset to save budget -# num_threads=1, -# **kwargs -# ) -# results[name] = float(best) -# except Exception as e: -# # If budget constraint causes early termination or other issues, record as 0 -# print(f"Algorithm {name} encountered error with budget constraint: {e}") -# results[name] = 0.0 - -# # Sanity check that we produced some results -# assert set(results.keys()) == {"GEPA-Base-Budget", "GEPA-UCB-Budget", "GEPA-Beam-Budget"} -# for v in results.values(): -# assert isinstance(v, float) -# assert v >= 0.0 # Should be non-negative scores - From 6b36c820696d4a4910cda0c0892922216a38961d Mon Sep 17 00:00:00 2001 From: doxav Date: Thu, 6 Nov 2025 19:44:12 +0100 Subject: [PATCH 5/6] Improved GEPA-UCB with score_range clamping --- opto/features/gepa/gepa_algorithms.py | 27 +- opto/trainer/algorithms/gepa_algorithms.py | 841 --------------------- 2 files changed, 23 insertions(+), 845 deletions(-) delete mode 100644 opto/trainer/algorithms/gepa_algorithms.py diff --git a/opto/features/gepa/gepa_algorithms.py b/opto/features/gepa/gepa_algorithms.py index 7494c0ca..5684aba2 100644 --- a/opto/features/gepa/gepa_algorithms.py +++ b/opto/features/gepa/gepa_algorithms.py @@ -224,6 +224,14 @@ class GEPAUCBSearch(UCBSearchAlgorithm): - Optional periodic Merge crossover (uniform per-parameter) with desirability checks """ + def _rank(self, raw: float) -> float: + """ + If a score_range is provided (lo, hi), clamp the scalar score into that band. + This keeps UCB-like behavior numerically stable without changing external APIs. + """ + if getattr(self, "score_range", None) is None or raw is None: return raw + lo, hi = self.score_range; return float(min(hi, max(lo, raw))) + def __init__(self, agent, optimizer=None, @@ -270,6 +278,7 @@ def train(self, pareto_subset_size: int = 24, num_search_iterations: int = 120, train_batch_size: int = 2, + score_range: Optional[Tuple[float, float]] = None, merge_every: int = 6, log_frequency: Optional[int] = None, save_frequency: Optional[int] = None, @@ -282,6 +291,7 @@ def train(self, num_threads = num_threads or self.num_threads log_frequency = log_frequency or 5 validate_ds = validation_dataset or train_dataset + self.score_range = score_range # Optional score clamping band for mean-based selections # Fix a Pareto subset (small, stable) to compute per-instance vectors assert len(validate_ds["inputs"]) > 0, "Empty dataset." @@ -296,10 +306,13 @@ def train(self, # Seed with current params base_params = {p: copy.deepcopy(p.data) for p in self.optimizer.parameters} - v0, m0 = self._evaluate_on_pareto(base_params, guide, num_threads=num_threads) - buffer.append(Candidate(params=base_params, eval_vector=v0, mean=m0, id=self._next_id(), ancestors=set())) + v0, m0_raw = self._evaluate_on_pareto(base_params, guide, num_threads=num_threads) + m0 = self._rank(m0_raw) + buffer.append(Candidate(params=base_params, eval_vector=v0, mean=m0, id=self._next_id(), + ancestors=set(), meta={"raw_mean": m0_raw})) print_color(f"[GEPA] Seed candidate mean={m0:.4f}", "cyan") + metrics = {"best_means": [], "new_child_means": [], "merge_accepts": 0, "total_merges": 0} for it in range(1, num_search_iterations + 1): @@ -322,14 +335,16 @@ def train(self, continue # Evaluate child on Pareto subset - child_vec, child_mean = self._evaluate_on_pareto(update_dict, guide, num_threads=num_threads) + child_vec, child_mean_raw = self._evaluate_on_pareto(update_dict, guide, num_threads=num_threads) + child_mean = self._rank(child_mean_raw) child = Candidate(params=update_dict, eval_vector=child_vec, mean=child_mean, id=self._next_id(), parent_ids=(parent.id,), ancestors=set(parent.ancestors) | {parent.id}, - created_iter=it) + created_iter=it, + meta={"raw_mean": child_mean_raw}) buffer.append(child) metrics["new_child_means"].append(child_mean) print_color(f"[GEPA] iter {it}: child mean={child_mean:.4f} (train-batch≈{train_batch_mean})", "green") @@ -347,6 +362,10 @@ def train(self, if merged is not None: merged.id = self._next_id() merged.created_iter = it + # preserve raw and clamp to range for ranking/logging + _raw = merged.mean + merged.meta["raw_mean"] = _raw + merged.mean = self._rank(_raw) buffer.append(merged) metrics["merge_accepts"] += 1 print_color(f"[GEPA] Merge accepted: mean={merged.mean:.4f}", "magenta") diff --git a/opto/trainer/algorithms/gepa_algorithms.py b/opto/trainer/algorithms/gepa_algorithms.py deleted file mode 100644 index c0283e38..00000000 --- a/opto/trainer/algorithms/gepa_algorithms.py +++ /dev/null @@ -1,841 +0,0 @@ -# opto/trainer/algorithms/gepa_algorithms.py -# GEPA (+Merge) algorithms for Trace -# - GEPAUCBSearch: subclass of UCBSearchAlgorithm -# - GEPABeamPareto: subclass of BeamsearchAlgorithm (Pareto select + single-parent incremental) -# - GEPATrainer: subclass of Trainer (minimal GEPA loop) -# -# All default to OptoPrimeV2 if optimizer=None. - -from __future__ import annotations -import copy -import math -import random -import functools -import types -from dataclasses import dataclass, field -from typing import Any, Dict, List, Optional, Tuple - -import numpy as np - -from opto.optimizers.optoprime_v2 import OptoPrimeV2 -from opto.trace.nodes import ParameterNode -from opto.trainer.algorithms.UCBsearch import UCBSearchAlgorithm -from opto.trainer.algorithms.beamsearch_algorithm import BeamsearchAlgorithm -from opto.trainer.algorithms.algorithm import Trainer -from opto.trainer.algorithms.basic_algorithms import ( - evaluate, - batchify, - standard_optimization_step, -) -from opto.trainer.utils import async_run -# Prefer thread-safe batched runner (deep-copies per task). Fallback handled at callsite. -try: - from opto.trainer.utils import batch_run # type: ignore -except Exception: # pragma: no cover - batch_run = None -from opto.optimizers.utils import print_color - - -# ----------------------------- Utilities ----------------------------- # - -@dataclass -class Candidate: - params: Dict[ParameterNode, Any] - eval_vector: List[float] # per-instance scores on fixed Pareto subset - mean: float - id: int - parent_ids: Tuple[int, ...] = field(default_factory=tuple) - ancestors: set = field(default_factory=set) - created_iter: int = 0 - wins: int = 0 # updated by Pareto accounting - meta: Dict[str, Any] = field(default_factory=dict) # freeform - -def _eval_on_subset(agent, guide, xs, infos, *, num_threads: Optional[int], desc: str) -> List[float]: - return evaluate(agent, guide, xs, infos, min_score=None, num_threads=num_threads, description=desc) - -def _compute_pareto_counts(cands: List[Candidate]) -> None: - """ - "Best-for-at-least-one-instance" winners. - For each position m in eval vectors, find argmax candidate and credit a win. - """ - if not cands: - return - L = len(cands[0].eval_vector) - # Reset - for c in cands: - c.wins = 0 - # Credit wins - for m in range(L): - best_idx = None - best_val = -float("inf") - for i, c in enumerate(cands): - v = c.eval_vector[m] if m < len(c.eval_vector) else -float("inf") - if v > best_val: - best_val, best_idx = v, i - if best_idx is not None: - cands[best_idx].wins += 1 - -def _pareto_sample(cands: List[Candidate], *, temperature: float = 1.0, rng: random.Random) -> Candidate: - """ - Sample a parent from union of per-instance winners, proportional to wins^1/T. - """ - if not cands: - raise ValueError("Empty candidate buffer.") - _compute_pareto_counts(cands) - wins = np.array([max(1, c.wins) for c in cands], dtype=float) # avoid zero - if temperature <= 0: - # Deterministic pick - return cands[int(wins.argmax())] - weights = wins ** (1.0 / max(1e-6, temperature)) - probs = weights / (weights.sum() if weights.sum() > 0 else 1.0) - idx = rng.choices(range(len(cands)), weights=probs, k=1)[0] - return cands[idx] - -def _uniform_merge_params(a: Dict[ParameterNode, Any], b: Dict[ParameterNode, Any], rng: random.Random) -> Dict[ParameterNode, Any]: - """ - Simple, robust "crossover": per-parameter uniform pick between parents. - (System-aware enough for prompt/code params, cheap, and safe.) - """ - keys = set(a.keys()) | set(b.keys()) - merged: Dict[ParameterNode, Any] = {} - for p in keys: - if p in a and p in b: - merged[p] = copy.deepcopy(a[p] if rng.random() < 0.5 else b[p]) - elif p in a: - merged[p] = copy.deepcopy(a[p]) - else: - merged[p] = copy.deepcopy(b[p]) - return merged - -def _maybe_merge(buffer: List[Candidate], - *, - agent, - guide, - pareto_inputs: List[Any], - pareto_infos: List[Any], - num_threads: Optional[int], - rng: random.Random, - tried_pairs: set, - max_tries: int = 8) -> Optional[Candidate]: - """ - Try merging two non-lineage candidates once; return merged if better than both parents' mean, else None. - """ - if len(buffer) < 2: - return None - # Prefer winners - _compute_pareto_counts(buffer) - pool = sorted(buffer, key=lambda c: (c.wins, c.mean), reverse=True) - - # Try a few distinct pairs - for _ in range(max_tries): - i, j = rng.sample(range(len(pool)), 2) - a, b = pool[i], pool[j] - if a.id == b.id: - continue - if a.id in b.ancestors or b.id in a.ancestors: - continue # avoid direct ancestry - key = tuple(sorted((a.id, b.id))) - if key in tried_pairs: - continue - tried_pairs.add(key) - - merged_params = _uniform_merge_params(a.params, b.params, rng) - # Evaluate merged on Pareto subset - original_params = _snapshot_params_fast(list(agent.parameters())) - try: - # load params to agent - from opto.optimizers.optimizer import Optimizer # type: ignore - # We only need the parameters dict projection; we can set via optimizer.update if available - # But we don't have an optimizer here; use ParameterNode._set - for p, v in merged_params.items(): - p._set(v) - - vec = _eval_on_subset(agent, guide, pareto_inputs, pareto_infos, num_threads=num_threads, - desc="GEPA+Merge: evaluating merged") - mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - finally: - # restore original - for p, v in original_params.items(): - p._set(v) - - if mean > max(a.mean, b.mean): - merged = Candidate(params=merged_params, - eval_vector=vec, - mean=mean, - id=-1, # to be set by caller - parent_ids=(a.id, b.id), - ancestors=set(a.ancestors) | set(b.ancestors) | {a.id, b.id}, - created_iter=0) - return merged - return None - - -def _maybe_merge_ancestor_aware( - buffer: List[Candidate], - *, - id2cand: Dict[int, Candidate], - module_groups: List[List[ParameterNode]], - agent, - guide, - optimizer, - train_dataset: Dict[str, List[Any]], - train_batch_size: int, - pareto_inputs: List[Any], - pareto_infos: List[Any], - num_threads: Optional[int], - rng: random.Random, - tried_pairs: set, - budget_tracker: Optional[Dict[str, int]] = None, - budget_B: Optional[int] = None, - max_tries: int = 8 -) -> Optional[Tuple[Candidate, int]]: - """ - Ancestor-aware merge with budget tracking. Returns (merged_candidate, rollouts_used). - """ - if len(buffer) < 2: - return None - - rollouts_used = 0 - - # Sample training minibatch (no replacement → lower variance) - k = min(train_batch_size, len(train_dataset["inputs"])) - idxs = np.random.choice(len(train_dataset["inputs"]), k, replace=False) - tx = [train_dataset["inputs"][i] for i in idxs] - ti = [train_dataset["infos"][i] for i in idxs] - - # Prefer winners for parent selection - _compute_pareto_counts(buffer) - pool = sorted(buffer, key=lambda c: (c.wins, c.mean), reverse=True) - - for _ in range(max_tries): - i, j = rng.sample(range(len(pool)), 2) - ci, cj = pool[i], pool[j] - if ci.id == cj.id: - continue - if ci.id in cj.ancestors or cj.id in ci.ancestors: - continue # avoid direct ancestry - key = tuple(sorted((ci.id, cj.id))) - if key in tried_pairs: - continue - tried_pairs.add(key) - - merged_params = _uniform_merge_params(ci.params, cj.params, rng) - - # Quick minibatch acceptability check - def _batch_mean_for(param_dict): - original = _snapshot_params_fast(list(optimizer.parameters)) - try: - _apply_params(optimizer, param_dict) - vec = evaluate(agent, guide, tx, ti, min_score=None, num_threads=num_threads, - description="MERGE(mini-batch accept)") - finally: - _apply_params(optimizer, original) - return float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - - rollouts_used += k - merged_batch_mean = _batch_mean_for(merged_params) - parent_means = [_batch_mean_for(ci.params), _batch_mean_for(cj.params)] - rollouts_used += 2 * k - - # Early budget guard (3*k minibatch evals) before Pareto eval - if budget_B is not None and budget_tracker is not None: - if budget_tracker["used"] + 3 * k + len(pareto_inputs) > budget_B: - return None - budget_tracker["used"] += 3 * k - - if merged_batch_mean <= max(parent_means): - continue # Not promising enough - - # Full Pareto evaluation - original = _snapshot_params_fast(list(optimizer.parameters)) - try: - _apply_params(optimizer, merged_params) - vec = evaluate(agent, guide, pareto_inputs, pareto_infos, min_score=None, - num_threads=num_threads, description="GEPA+Merge: ancestor-aware Pareto eval") - finally: - _apply_params(optimizer, original) - mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - # Account Pareto evaluation cost in the global budget and local counter. - if budget_B is not None and budget_tracker is not None: - budget_tracker["used"] += len(pareto_inputs) - rollouts_used += len(pareto_inputs) - - merged = Candidate(params=merged_params, - eval_vector=vec, mean=mean, - id=-1, parent_ids=(ci.id, cj.id), - ancestors=set(ci.ancestors) | set(cj.ancestors) | {ci.id, cj.id}, - created_iter=0) - return merged, rollouts_used - - return None - - -def _ensure_optimizer(agent, optimizer): - if optimizer is not None: - return optimizer - params = [p for p in agent.parameters()] # List[ParameterNode] - return OptoPrimeV2(parameters=params) - - -def _train_step_generate_child(agent, guide, optimizer, train_xs, train_infos, *, verbose=False, num_threads=None): - """ - Single-parent, incremental evolution "mutation": run forward on a minibatch to get batched feedback, - then optimizer.step(bypassing=True) to obtain a new candidate param dict (without applying). - """ - use_parallel = (num_threads is not None and num_threads > 1 and batch_run is not None) - if use_parallel: - # Pre-bind args → pass callables only. Robust to different batch_run signatures. - callables = [ - functools.partial(standard_optimization_step, agent, x, guide, info) - for x, info in zip(train_xs, train_infos) - ] - try: - outputs = batch_run( - callables, - max_workers=num_threads, - description="GEPA forward (mutate parent)", - ) - except TypeError: - # Fallback: older/other signature (e.g., batch_run(callables, max_workers)) - try: - outputs = batch_run(callables, num_threads) - except Exception: - outputs = None - # Normalize outputs to a list of results. batch_run in different versions may: - # - return the list of results, - # - return a callable that returns the results, - # - return a generator/iterator, - # - or return None. - try: - if callable(outputs): - outputs = outputs() - elif isinstance(outputs, types.GeneratorType): - outputs = list(outputs) - elif outputs is None: - # fallback to sequential evaluation - outputs = [fn() for fn in callables] - elif not isinstance(outputs, (list, tuple)): - # Some other iterable (e.g. map object) - try: - outputs = list(outputs) - except Exception: - outputs = [fn() for fn in callables] - except Exception: - # Any error while normalizing → fallback to sequential - outputs = [fn() for fn in callables] - else: - # Safe sequential fallback. - outputs = [standard_optimization_step(agent, x, guide, info) for x, info in zip(train_xs, train_infos)] - - scores, targets, feedbacks = [], [], [] - for target, score, feedback in outputs: - scores.append(score) - targets.append(target) - feedbacks.append(feedback) - - target_batch = batchify(*targets) - feedback_batch = batchify(*feedbacks).data - - optimizer.zero_feedback() - optimizer.backward(target_batch, feedback_batch) - try: - update_dict = optimizer.step(bypassing=True, verbose=("output" if verbose else False)) - if not isinstance(update_dict, dict) or len(update_dict) == 0: - # Fallback: treat current as child (rare) - update_dict = {p: copy.deepcopy(p.data) for p in optimizer.parameters} - except Exception as e: - print_color(f"[GEPA] optimizer.step error: {e}", "red") - update_dict = {} - return update_dict, (None if not scores or any(s is None for s in scores) else float(np.mean(scores))) - - -def _apply_params(optimizer, param_dict: Dict[ParameterNode, Any]): - """Load param dict into the agent via optimizer.update (preserves projections).""" - optimizer.update(param_dict) - - -def _snapshot_params_fast(parameters: List[ParameterNode]) -> Dict[ParameterNode, Any]: - """ - Snapshot ParameterNode->value with minimal copying: - - immutables (str/int/float/bool/tuple/bytes/None): no copy - - numpy arrays: .copy() - - everything else: deepcopy (safe fallback) - """ - snap: Dict[ParameterNode, Any] = {} - immutables = (str, int, float, bool, tuple, frozenset, bytes, type(None)) - for p in parameters: - v = getattr(p, "data", None) - if isinstance(v, immutables): - snap[p] = v - elif isinstance(v, np.ndarray): - snap[p] = v.copy() - else: - snap[p] = copy.deepcopy(v) - return snap - - -def _fingerprint_params(params_dict: Dict[ParameterNode, Any]) -> Tuple: - """ - Hashable fingerprint of a ParameterNode->value dict for optional caching. - Uses (param-id, repr(value)) with special handling for numpy arrays. - """ - items: List[Tuple] = [] - for p, v in params_dict.items(): - pid = getattr(p, "uid", None) or getattr(p, "name", None) or id(p) - try: - if isinstance(v, np.ndarray): - items.append(("arr", pid, v.shape, v.dtype.str, hash(v.tobytes()))) - else: - items.append(("val", pid, repr(v))) - except Exception: - items.append(("val", pid, repr(v))) - return tuple(sorted(items)) - - -# ======================= Variant 1: GEPA + Merge (UCB subclass) ======================= # - -class GEPAUCBSearch(UCBSearchAlgorithm): - """ - GEPA (+Merge) implemented atop UCBSearchAlgorithm. - Differences vs base UCB: - - Fixed Pareto subset (D_pareto) and per-instance vectors kept for each candidate - - Parent selection = Pareto "best-for-at-least-one" sampling (wins-weighted); UCB used only for eviction fallback - - Single-parent incremental mutation via a minibatch - - Optional periodic Merge crossover (uniform per-parameter) with desirability checks - """ - - def __init__(self, - agent, - optimizer=None, - *, - max_buffer_size: int = 16, - ucb_exploration_factor: float = 0.8, - rng_seed: int = 7, - logger=None, - num_threads: Optional[int] = None, - module_groups: Optional[Dict[str, List[ParameterNode]] | List[List[ParameterNode]]] = None, - selectmodule_policy: str = "round_robin", - enable_pareto_cache: bool = False): - optimizer = _ensure_optimizer(agent, optimizer) - super().__init__(agent, optimizer, - max_buffer_size=max_buffer_size, - ucb_exploration_factor=ucb_exploration_factor, - logger=logger, - num_threads=num_threads) - self.rng = random.Random(rng_seed) - np.random.seed(rng_seed) # ensure numpy reproducibility for np.random.choice - self._pareto_inputs: List[Any] = [] - self._pareto_infos: List[Any] = [] - self._id_counter = 0 - self.enable_pareto_cache = enable_pareto_cache - self._pareto_cache: Dict[Tuple, Tuple[List[float], float]] = {} - - def _next_id(self) -> int: - self._id_counter += 1 - return self._id_counter - - def _evaluate_on_pareto(self, params_dict: Dict[ParameterNode, Any], guide, *, num_threads) -> Tuple[List[float], float]: - cache_key = _fingerprint_params(params_dict) if self.enable_pareto_cache else None - if cache_key is not None: - cached = self._pareto_cache.get(cache_key) - if cached is not None: - return cached - original_params = _snapshot_params_fast(list(self.optimizer.parameters)) - try: - _apply_params(self.optimizer, params_dict) - vec = _eval_on_subset(self.agent, guide, self._pareto_inputs, self._pareto_infos, - num_threads=num_threads, desc="GEPA: evaluate on Pareto subset") - mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - result = (vec, mean) - if cache_key is not None: - self._pareto_cache[cache_key] = result - return result - finally: - _apply_params(self.optimizer, original_params) - - def _select_pareto_parent(self, cand_buffer: List[Candidate]) -> Candidate: - return _pareto_sample(cand_buffer, temperature=1.0, rng=self.rng) - - def train(self, - guide, - train_dataset: Dict[str, List[Any]], - *, - validate_dataset: Optional[Dict[str, List[Any]]] = None, - pareto_subset_size: int = 24, - num_search_iterations: int = 120, - train_batch_size: int = 2, - merge_every: int = 6, - log_frequency: Optional[int] = None, - save_frequency: Optional[int] = None, - save_path: str = "checkpoints/gepa_ucb_agent.pkl", - verbose: bool = False, - num_threads: Optional[int] = None) -> Tuple[Dict[str, Any], float]: - """ - GEPA search loop with Pareto sampling + (optional) Merge. - """ - num_threads = num_threads or self.num_threads - log_frequency = log_frequency or 5 - validate_ds = validate_dataset or train_dataset - - # Fix a Pareto subset (small, stable) to compute per-instance vectors - assert len(validate_ds["inputs"]) > 0, "Empty dataset." - idxs = np.random.choice(len(validate_ds["inputs"]), - min(pareto_subset_size, len(validate_ds["inputs"])), - replace=False) - self._pareto_inputs = [validate_ds["inputs"][i] for i in idxs] - self._pareto_infos = [validate_ds["infos"][i] for i in idxs] - - buffer: List[Candidate] = [] - tried_merges: set = set() - id2cand: Dict[int, Candidate] = {} - - # Seed with current params - base_params = _snapshot_params_fast(list(self.optimizer.parameters)) - v0, m0 = self._evaluate_on_pareto(base_params, guide, num_threads=num_threads) - seed = Candidate(params=base_params, eval_vector=v0, mean=m0, id=self._next_id(), ancestors=set(), created_iter=0) - buffer.append(seed) - id2cand[seed.id] = seed - print_color(f"[GEPA] Seed candidate mean={m0:.4f}", "cyan") - - metrics = {"best_means": [], "new_child_means": [], "merge_accepts": 0, "total_merges": 0} - - for it in range(1, num_search_iterations + 1): - # Select parent by Pareto winners - parent = self._select_pareto_parent(buffer) - _apply_params(self.optimizer, parent.params) - - # Sample train minibatch - train_size = min(train_batch_size, len(train_dataset["inputs"])) - tr_idxs = np.random.choice(len(train_dataset["inputs"]), train_size, replace=False) - train_xs = [train_dataset["inputs"][i] for i in tr_idxs] - train_info = [train_dataset["infos"][i] for i in tr_idxs] - - # Generate child via one incremental step - update_dict, train_batch_mean = _train_step_generate_child( - self.agent, guide, self.optimizer, train_xs, train_info, verbose=verbose, num_threads=num_threads - ) - if not update_dict: - print_color("[GEPA] Empty child update; skipping.", "yellow") - continue - - # Evaluate child on Pareto subset - child_vec, child_mean = self._evaluate_on_pareto(update_dict, guide, num_threads=num_threads) - child = Candidate(params=update_dict, - eval_vector=child_vec, - mean=child_mean, - id=self._next_id(), - parent_ids=(parent.id,), - ancestors=set(parent.ancestors) | {parent.id}, - created_iter=it) - buffer.append(child) - metrics["new_child_means"].append(child_mean) - print_color(f"[GEPA] iter {it}: child mean={child_mean:.4f} (train-batch≈{train_batch_mean})", "green") - - # Optional Merge - if merge_every and (it % merge_every == 0): - metrics["total_merges"] += 1 - merged = _maybe_merge(buffer, - agent=self.agent, guide=guide, - pareto_inputs=self._pareto_inputs, - pareto_infos=self._pareto_infos, - num_threads=num_threads, - rng=self.rng, - tried_pairs=tried_merges) - if merged is not None: - merged.id = self._next_id() - merged.created_iter = it - buffer.append(merged) - metrics["merge_accepts"] += 1 - print_color(f"[GEPA] Merge accepted: mean={merged.mean:.4f}", "magenta") - - # Keep buffer bounded: remove the candidate with lowest (wins, mean) - if len(buffer) > self.max_buffer_size: - _compute_pareto_counts(buffer) - buffer.sort(key=lambda c: (c.wins, c.mean)) - evicted = buffer.pop(0) - print_color(f"[GEPA] Evicted cand#{evicted.id} (wins={evicted.wins}, mean={evicted.mean:.4f})", "yellow") - - # Track & log - best = max(buffer, key=lambda c: c.mean) - metrics["best_means"].append(best.mean) - if it % log_frequency == 0: - self.logger.log("GEPA best mean", best.mean, it, color="green") - - # Save best candidate snapshot (optional) - if save_frequency and it % save_frequency == 0: - _apply_params(self.optimizer, best.params) - self.save_agent(save_path, it) - - # Load best into the agent and return - best = max(buffer, key=lambda c: c.mean) if buffer else buffer[0] - _apply_params(self.optimizer, best.params) - return metrics, float(best.mean) - - -# ================= Variant 2: Beamsearch subclass with Pareto select ================= # - -class GEPABeamPareto(BeamsearchAlgorithm): - """ - BeamsearchAlgorithm retrofit: - - override select() to a Pareto "best-for-at-least-one" selector - - replace deep beam expansion with GEPA’s single-parent incremental evolution - """ - - def __init__(self, agent, optimizer=None, *, rng_seed: int = 11, logger=None, - num_threads: Optional[int] = None, - module_groups: Optional[Dict[str, List[ParameterNode]] | List[List[ParameterNode]]] = None, - selectmodule_policy: str = "round_robin"): - optimizer = _ensure_optimizer(agent, optimizer) - super().__init__(agent, optimizer, num_threads=num_threads, logger=logger) - self.rng = random.Random(rng_seed) - np.random.seed(rng_seed) - - # We keep a Pareto select helper that returns (selected_params, wins, scores) - def select(self, - candidates: List[Dict[ParameterNode, Any]], - validate_guide, - validation_mini_dataset, - beam_width: int, - num_threads: int = None, - min_score: float = None, - return_scores: bool = False): - """ - Override to Pareto union-of-winners on the mini validation batch. - """ - # Evaluate each candidate to a vector on the mini validation - cand_objs: List[Candidate] = [] - current_params = _snapshot_params_fast(list(self.optimizer.parameters)) - try: - for idx, params in enumerate(candidates): - _apply_params(self.optimizer, params) - vec = evaluate(self.agent, - validate_guide, - validation_mini_dataset['inputs'], - validation_mini_dataset['infos'], - min_score=min_score, - num_threads=num_threads, - description=f"Validating candidate {idx+1}/{len(candidates)} (Pareto)") - mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - cand_objs.append(Candidate(params=params, eval_vector=vec, mean=mean, id=idx)) - finally: - _apply_params(self.optimizer, current_params) - - # Compute wins and select top "beam_width" by (wins, mean) - _compute_pareto_counts(cand_objs) - cand_objs.sort(key=lambda c: (c.wins, c.mean), reverse=True) - selected = cand_objs[: min(beam_width, len(cand_objs))] - sel_params = [c.params for c in selected] - sel_scores = [c.mean for c in selected] - if return_scores: - return sel_params, sel_scores - return sel_params - - # Replace beam "train" with GEPA-style incremental loop (keeps BeamsearchAlgorithm API) - def train(self, guide, train_dataset, *, - validate_dataset=None, pareto_subset_size: int = 24, - num_search_iterations: int = 120, train_batch_size: int = 2, - merge_every: int = 6, log_frequency: Optional[int] = None, - save_frequency: Optional[int] = None, - save_path: str = "checkpoints/gepa_beam_agent.pkl", - verbose: bool = False, num_threads: Optional[int] = None, - module_groups: Optional[Dict[str, List[ParameterNode]] | List[List[ParameterNode]]] = None, - selectmodule_policy: str = "round_robin", - budget_B: Optional[int] = None, - accept_epsilon: float = 0.0): - num_threads = num_threads or self.num_threads - log_frequency = log_frequency or 5 - validate_ds = validate_dataset or train_dataset - - # Fix Pareto subset for this run - idxs = np.random.choice(len(validate_ds["inputs"]), - min(pareto_subset_size, len(validate_ds["inputs"])), - replace=False) - pareto_inputs = [validate_ds["inputs"][i] for i in idxs] - pareto_infos = [validate_ds["infos"][i] for i in idxs] - - # Seed buffer - buffer: List[Candidate] = [] - base_params = _snapshot_params_fast(list(self.optimizer.parameters)) - original = _snapshot_params_fast(list(self.optimizer.parameters)) - try: - _apply_params(self.optimizer, base_params) - vec = evaluate(self.agent, guide, pareto_inputs, pareto_infos, - min_score=None, num_threads=num_threads, - description="GEPA(beam): seed evaluation") - finally: - _apply_params(self.optimizer, original) - m0 = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - buffer.append(Candidate(params=base_params, eval_vector=vec, mean=m0, id=0, ancestors=set())) - tried_merges: set = set() - - best_mean = m0 - for it in range(1, num_search_iterations + 1): - # Pareto-select parent and mutate - _compute_pareto_counts(buffer) - parent = _pareto_sample(buffer, temperature=1.0, rng=self.rng) - _apply_params(self.optimizer, parent.params) - - # Make a child - k = min(train_batch_size, len(train_dataset["inputs"])) - tr = np.random.choice(len(train_dataset["inputs"]), k, replace=False) - train_xs = [train_dataset["inputs"][i] for i in tr] - train_in = [train_dataset["infos"][i] for i in tr] - - update_dict, _ = _train_step_generate_child(self.agent, guide, self.optimizer, train_xs, train_in, - verbose=verbose, num_threads=num_threads) - if not update_dict: - continue - - # Evaluate child on Pareto subset - original = _snapshot_params_fast(list(self.optimizer.parameters)) - try: - _apply_params(self.optimizer, update_dict) - vec = evaluate(self.agent, guide, pareto_inputs, pareto_infos, min_score=None, - num_threads=num_threads, description="GEPA(beam): child eval") - finally: - _apply_params(self.optimizer, original) - mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - buffer.append(Candidate(params=update_dict, eval_vector=vec, mean=mean, id=len(buffer), - parent_ids=(parent.id,), ancestors=set(parent.ancestors) | {parent.id})) - best_mean = max(best_mean, mean) - if it % log_frequency == 0: - self.logger.log("GEPA(beam) best mean", best_mean, it, color="green") - - # Periodic merge - if merge_every and it % merge_every == 0: - merged = _maybe_merge(buffer, - agent=self.agent, guide=guide, - pareto_inputs=pareto_inputs, pareto_infos=pareto_infos, - num_threads=num_threads, rng=self.rng, tried_pairs=tried_merges) - if merged is not None: - merged.id = len(buffer) - buffer.append(merged) - - # Trim buffer softly (keep top by (wins, mean)) - if len(buffer) > 16: - _compute_pareto_counts(buffer) - buffer.sort(key=lambda c: (c.wins, c.mean), reverse=True) - buffer[:] = buffer[:16] - - # Optional save - if save_frequency and it % save_frequency == 0: - best = max(buffer, key=lambda c: c.mean) - _apply_params(self.optimizer, best.params) - self.save_agent(save_path, it) - - best = max(buffer, key=lambda c: c.mean) - _apply_params(self.optimizer, best.params) - return {"best_mean": best.mean}, float(best.mean) - - -# =================== Variant 3: Minimal GEPA on AlgorithmBase =================== # - -class GEPAAlgorithmBase(Trainer): - """ - Lightweight GEPA (+Merge) with only Trainer dependency. - Useful when you want the simplest control loop with your own logging/saving. - """ - - def __init__(self, agent, optimizer=None, *, rng_seed: int = 13, logger=None, - num_threads: Optional[int] = None, - module_groups: Optional[Dict[str, List[ParameterNode]] | List[List[ParameterNode]]] = None, - selectmodule_policy: str = "round_robin"): - super().__init__(agent, num_threads=num_threads, logger=logger) - self.optimizer = _ensure_optimizer(agent, optimizer) - self.rng = random.Random(rng_seed) - np.random.seed(rng_seed) - - def train(self, - guide, - train_dataset, - *, - validate_dataset=None, - pareto_subset_size: int = 24, - num_iters: int = 100, - train_batch_size: int = 2, - merge_every: int = 5, - num_threads: Optional[int] = None, - save_path: Optional[str] = None): - num_threads = num_threads or self.num_threads - validate_ds = validate_dataset or train_dataset - - # Pareto subset - idxs = np.random.choice(len(validate_ds["inputs"]), - min(pareto_subset_size, len(validate_ds["inputs"])), - replace=False) - xsP = [validate_ds["inputs"][i] for i in idxs] - isP = [validate_ds["infos"][i] for i in idxs] - - # Seed - buffer: List[Candidate] = [] - base_params = _snapshot_params_fast(list(self.optimizer.parameters)) - original = _snapshot_params_fast(list(self.optimizer.parameters)) - try: - _apply_params(self.optimizer, base_params) - vec = evaluate(self.agent, guide, xsP, isP, min_score=None, num_threads=num_threads, - description="GEPA(base): seed eval") - finally: - _apply_params(self.optimizer, original) - m0 = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - buffer.append(Candidate(params=base_params, eval_vector=vec, mean=m0, id=0, ancestors=set())) - tried_merges: set = set() - - for it in range(1, num_iters + 1): - # Parent select - _compute_pareto_counts(buffer) - parent = _pareto_sample(buffer, temperature=1.0, rng=self.rng) - _apply_params(self.optimizer, parent.params) - - # Child - k = min(train_batch_size, len(train_dataset["inputs"])) - tr = np.random.choice(len(train_dataset["inputs"]), k, replace=False) - tx = [train_dataset["inputs"][i] for i in tr] - ti = [train_dataset["infos"][i] for i in tr] - update_dict, _ = _train_step_generate_child(self.agent, guide, self.optimizer, tx, ti, - verbose=False, num_threads=num_threads) - if not update_dict: - continue - - # Eval child - original = _snapshot_params_fast(list(self.optimizer.parameters)) - try: - _apply_params(self.optimizer, update_dict) - vec = evaluate(self.agent, guide, xsP, isP, min_score=None, num_threads=num_threads, - description="GEPA(base): child eval") - finally: - _apply_params(self.optimizer, original) - mean = float(np.mean(vec)) if all(s is not None for s in vec) else -float("inf") - buffer.append(Candidate(params=update_dict, eval_vector=vec, mean=mean, id=len(buffer), - parent_ids=(parent.id,), ancestors=set(parent.ancestors) | {parent.id})) - - # Merge - if merge_every and it % merge_every == 0: - merged = _maybe_merge(buffer, - agent=self.agent, guide=guide, - pareto_inputs=xsP, pareto_infos=isP, - num_threads=num_threads, rng=self.rng, tried_pairs=tried_merges) - if merged is not None: - merged.id = len(buffer) - buffer.append(merged) - - # Keep compact buffer - if len(buffer) > 16: - _compute_pareto_counts(buffer) - buffer.sort(key=lambda c: (c.wins, c.mean), reverse=True) - buffer[:] = buffer[:16] - - # Log - best = max(buffer, key=lambda c: c.mean) - if self.logger: - self.logger.log("GEPA(base) best mean", best.mean, it, color="green") - - # Optional save - if save_path and it % 10 == 0: - _apply_params(self.optimizer, best.params) - self.save_agent(save_path, it) - - # Load best into agent - best = max(buffer, key=lambda c: c.mean) - _apply_params(self.optimizer, best.params) - return {"best_mean": best.mean}, float(best.mean) - From cb203b1dbd83918085ebd8278197b4164622c10a Mon Sep 17 00:00:00 2001 From: doxav Date: Thu, 6 Nov 2025 19:54:14 +0100 Subject: [PATCH 6/6] removed from test GEPA unavailable feature > GEPA-UCB might be the one to keep --- tests/unit_tests/test_gepa_algorithms.py | 217 +---------------------- 1 file changed, 3 insertions(+), 214 deletions(-) diff --git a/tests/unit_tests/test_gepa_algorithms.py b/tests/unit_tests/test_gepa_algorithms.py index c435628e..fdaa6c0b 100644 --- a/tests/unit_tests/test_gepa_algorithms.py +++ b/tests/unit_tests/test_gepa_algorithms.py @@ -16,14 +16,15 @@ from opto.trace.nodes import node as trace_node from opto.optimizers.optoprime_v2 import OptoPrimeV2 import pytest -from opto.trainer.algorithms.gepa_algorithms import ( +from opto.features.gepa.gepa_algorithms import ( GEPAAlgorithmBase, GEPAUCBSearch, GEPABeamPareto, _compute_pareto_counts, _pareto_sample, - _uniform_merge_params, + _uniform_merge_params ) +from opto.trainer.algorithms.beamsearch_algorithm import BeamsearchAlgorithm from opto.trainer.evaluators import evaluate from opto.trainer.guide import Guide from opto.utils.llm import DummyLLM @@ -166,126 +167,6 @@ def test_gepa_variants_converge_on_dummyllm(algo_cls, train_kwargs): assert agent.param.data == target_add -def test_compare_gepa_vs_basicsearch_on_dummyllm(): - from opto.trainer.algorithms.basic_algorithms import BasicSearchAlgorithm - - target_add = 7 - ds = make_dataset(target_add, n=6) - agent_gepa = AddAgent(param=0) - agent_basic = AddAgent(param=0) - - opt_gepa = build_optimizer(agent_gepa, suggest_value=target_add) - opt_basic = build_optimizer(agent_basic, suggest_value=target_add) - - # GEPA - gepa = GEPAAlgorithmBase(agent_gepa, optimizer=opt_gepa, logger=None, num_threads=1) - _, best_gepa = gepa.train( - guide=ExactMatchGuide(), - train_dataset=ds, - validate_dataset=ds, - pareto_subset_size=4, - num_iters=8, - train_batch_size=2, - merge_every=2, - num_threads=1, - ) - - # BasicSearch baseline - basic = BasicSearchAlgorithm(agent_basic, optimizer=opt_basic, logger=None, num_threads=1) - basic.train( - guide=ExactMatchGuide(), - train_dataset=ds, - validate_dataset=ds, - num_proposals=1, - num_epochs=1, - batch_size=1, - test_dataset=ds, - eval_frequency=1, - num_threads=1, - verbose=False, - ) - - # Evaluate both on full dataset - score_gepa = np.mean(evaluate(agent_gepa, ExactMatchGuide(), ds["inputs"], ds["infos"], num_threads=2)) - score_basic = np.mean(evaluate(agent_basic, ExactMatchGuide(), ds["inputs"], ds["infos"], num_threads=2)) - - assert best_gepa == pytest.approx(1.0, rel=0, abs=1e-6) - assert score_gepa == pytest.approx(1.0, rel=0, abs=1e-6) - assert score_basic == pytest.approx(1.0, rel=0, abs=1e-6) - - -def test_snapshot_params_fast(): - """Test the fast parameter snapshot utility function.""" - from opto.trainer.algorithms.gepa_algorithms import _snapshot_params_fast - - @trace_model - class MultiTypeAgent: - def __init__(self): - self.int_param = trace_node(42, trainable=True) - self.str_param = trace_node("hello", trainable=True) - self.float_param = trace_node(3.14, trainable=True) - self.list_param = trace_node([1, 2, 3], trainable=True) - self.dict_param = trace_node({"key": "value"}, trainable=True) - # Test numpy array - self.np_param = trace_node(np.array([1, 2, 3]), trainable=True) - - def forward(self, x): - return x + self.int_param - - agent = MultiTypeAgent() - params = list(agent.parameters()) - - # Test snapshot - snapshot = _snapshot_params_fast(params) - - # Check that all parameters are included - assert len(snapshot) == len(params) - - # Modify original values - agent.int_param._set(100) - agent.str_param._set("modified") - agent.np_param._set(np.array([4, 5, 6])) - - # Verify snapshot preserved original values - for p in params: - if p.py_name == "int_param": - assert snapshot[p] == 42 - elif p.py_name == "str_param": - assert snapshot[p] == "hello" - elif p.py_name == "np_param": - assert np.array_equal(snapshot[p], np.array([1, 2, 3])) - - -def test_fingerprint_params(): - """Test the parameter fingerprinting utility function.""" - from opto.trainer.algorithms.gepa_algorithms import _fingerprint_params - - @trace_model - class SimpleAgent: - def __init__(self): - self.a = trace_node(1, trainable=True) - self.b = trace_node("test", trainable=True) - - def forward(self, x): - return x + self.a - - agent = SimpleAgent() - params_dict = {p: p.data for p in agent.parameters()} - - # Test fingerprinting - fp1 = _fingerprint_params(params_dict) - fp2 = _fingerprint_params(params_dict) - - # Same parameters should produce same fingerprint - assert fp1 == fp2 - - # Different parameters should produce different fingerprint - agent.a._set(2) - params_dict2 = {p: p.data for p in agent.parameters()} - fp3 = _fingerprint_params(params_dict2) - assert fp1 != fp3 - - def test_numpy_seeding_reproducibility(): """Test that numpy seeding ensures reproducible behavior.""" target_add = 3 @@ -335,64 +216,6 @@ def test_numpy_seeding_reproducibility(): assert best_diff == pytest.approx(1.0, rel=0, abs=1e-6) -def test_gepa_ucb_pareto_cache(): - """Test Pareto cache functionality in GEPAUCBSearch.""" - target_add = 4 - ds = make_dataset(target_add, n=3) - agent = AddAgent(param=0) - optimizer = build_optimizer(agent, suggest_value=target_add) - - # Test with cache enabled - algo = GEPAUCBSearch(agent=agent, optimizer=optimizer, logger=None, num_threads=1, enable_pareto_cache=True) - - metrics, best = algo.train( - guide=ExactMatchGuide(), - train_dataset=ds, - validate_dataset=ds, - pareto_subset_size=2, - num_search_iterations=2, - train_batch_size=1, - merge_every=2, - num_threads=1, - ) - - # Should converge to perfect solution - assert best == pytest.approx(1.0, rel=0, abs=1e-6) - assert agent.param.data == target_add - - # Test that cache was used (should have some entries) - # Note: exact cache size depends on algorithm behavior, but should be non-empty if enabled - if hasattr(algo, '_pareto_cache'): - assert isinstance(algo._pareto_cache, dict) - - -def test_budget_tracking_functionality(): - """Test budget tracking in GEPA algorithms.""" - target_add = 2 - ds = make_dataset(target_add, n=4) - agent = AddAgent(param=0) - optimizer = build_optimizer(agent, suggest_value=target_add) - - # Test GEPABeamPareto with budget - algo = GEPABeamPareto(agent=agent, optimizer=optimizer, logger=None, num_threads=1) - - metrics, best = algo.train( - guide=ExactMatchGuide(), - train_dataset=ds, - validate_dataset=ds, - pareto_subset_size=3, - num_search_iterations=2, - train_batch_size=1, - merge_every=2, - budget_B=10, # Low budget to test tracking - num_threads=1, - ) - - # Should still achieve good results even with budget constraint - assert isinstance(best, float) - assert best >= 0.0 # Should be non-negative score - - def test_thread_safety_with_sequential_fallback(): """Test that algorithms work correctly with sequential fallback when batch_run unavailable.""" target_add = 1 @@ -435,37 +258,3 @@ def test_thread_safety_with_sequential_fallback(): assert best2 == pytest.approx(1.0, rel=0, abs=1e-6) assert agent2.param.data == target_add - -def test_gepa_ucb_selectmodule_policy(): - """Test different module selection policies in GEPAUCBSearch.""" - target_add = 6 - ds = make_dataset(target_add, n=3) - - # Test different selection policies - policies = ["round_robin"] # Could test more if other policies are available - - for policy in policies: - agent = AddAgent(param=0) - optimizer = build_optimizer(agent, suggest_value=target_add) - - algo = GEPAUCBSearch( - agent=agent, - optimizer=optimizer, - logger=None, - num_threads=1, - selectmodule_policy=policy - ) - - metrics, best = algo.train( - guide=ExactMatchGuide(), - train_dataset=ds, - validate_dataset=ds, - pareto_subset_size=2, - num_search_iterations=2, - train_batch_size=1, - merge_every=2, - num_threads=1, - ) - - assert best == pytest.approx(1.0, rel=0, abs=1e-6) - assert agent.param.data == target_add