827 lines
33 KiB
Python
827 lines
33 KiB
Python
"""Deep 编译流水线编排。"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import shutil
|
|
import sys
|
|
import tempfile
|
|
import uuid
|
|
from dataclasses import asdict, replace
|
|
from pathlib import Path
|
|
from typing import Any
|
|
|
|
from scripts.dynamic_compile.fast.models import RolloutTrace
|
|
from scripts.dynamic_compile.fast.storage import (
|
|
atomic_write_json,
|
|
atomic_write_jsonl,
|
|
atomic_write_text,
|
|
load_json,
|
|
package_hash,
|
|
read_jsonl,
|
|
sha256_file,
|
|
sha256_text,
|
|
)
|
|
from scripts.dynamic_compile.fast.optimization.analyzer import SemanticClient
|
|
|
|
from .adapters.benchflow import (
|
|
BenchFlowInput,
|
|
SkillsBenchDevelopmentRolloutRunner,
|
|
load_benchflow_input,
|
|
)
|
|
from .adapters.semantic import DeepAnalyzer, valid_evidence
|
|
from .core.models import (
|
|
CellScore,
|
|
Coordinate,
|
|
DIMENSIONS,
|
|
LocalEdit,
|
|
ScoreMatrix,
|
|
SkillUnit,
|
|
)
|
|
from .core.markdown import (
|
|
parse_paragraphs,
|
|
parse_sections,
|
|
preserve_unit_boundary,
|
|
replace_unit_text,
|
|
validate_edit,
|
|
)
|
|
|
|
|
|
ROLLOUTS = 3
|
|
MAX_SECTION_ITERATIONS = 6
|
|
MAX_PARAGRAPH_ITERATIONS = 3
|
|
MAX_EDIT_GENERATION_ATTEMPTS = 3
|
|
GAP_THRESHOLD = 0.375
|
|
REJECTION_LIMIT = 2
|
|
MAX_PARALLEL = 3
|
|
DEFAULT_MODEL = "ali/deepseek-v4-pro-0813"
|
|
|
|
|
|
def _log(message: str) -> None:
|
|
print(f"[deep] {message}", file=sys.stderr, flush=True)
|
|
|
|
|
|
def _trace_from_dict(value: dict[str, Any]) -> RolloutTrace:
|
|
return RolloutTrace(**value)
|
|
|
|
|
|
def _column_to_dict(column: dict[str, CellScore]) -> dict[str, Any]:
|
|
return {unit_id: asdict(cell) for unit_id, cell in column.items()}
|
|
|
|
|
|
def _column_from_dict(
|
|
value: dict[str, Any], traces: list[RolloutTrace]
|
|
) -> dict[str, CellScore]:
|
|
return {
|
|
unit_id: CellScore(
|
|
float(cell["score"]), valid_evidence(cell.get("evidence"), traces), str(cell["reason"])
|
|
)
|
|
for unit_id, cell in value.items()
|
|
}
|
|
|
|
|
|
class DeepLoop:
|
|
def __init__(
|
|
self,
|
|
run_dir: Path,
|
|
analyzer: DeepAnalyzer | None = None,
|
|
runner: Any | None = None,
|
|
):
|
|
self.run_dir = run_dir.resolve()
|
|
self.state_path = self.run_dir / "run.json"
|
|
state = load_json(self.state_path)
|
|
if not isinstance(state, dict):
|
|
raise ValueError(f"invalid or missing run state: {self.state_path}")
|
|
self.state = state
|
|
self.temp = self.run_dir / ".tmp"
|
|
self.current = self.temp / "current"
|
|
self.model = state.get("model", state.get("semantic_model", DEFAULT_MODEL))
|
|
self.analyzer = analyzer or DeepAnalyzer(
|
|
SemanticClient(self.model), MAX_PARALLEL,
|
|
)
|
|
context = BenchFlowInput(
|
|
state["task"]["name"],
|
|
Path(state["task"]["directory"]),
|
|
state["task"]["agent"],
|
|
state["task"]["model"],
|
|
state["task"]["prompt"],
|
|
[],
|
|
)
|
|
self.runner = runner or SkillsBenchDevelopmentRolloutRunner(
|
|
context,
|
|
Path(tempfile.gettempdir()) / "skill-compiler-deep" / state["run_id"],
|
|
MAX_PARALLEL,
|
|
archive_root=self.run_dir / "rollouts",
|
|
)
|
|
|
|
@classmethod
|
|
def create(
|
|
cls,
|
|
skill: Path,
|
|
traces: Path,
|
|
output: Path | None = None,
|
|
*,
|
|
model: str = DEFAULT_MODEL,
|
|
analyzer: DeepAnalyzer | None = None,
|
|
runner: Any | None = None,
|
|
) -> "DeepLoop":
|
|
skill = skill.resolve()
|
|
if not (skill / "SKILL.md").is_file():
|
|
raise ValueError("--skill must be a skill package containing SKILL.md")
|
|
context = load_benchflow_input(traces)
|
|
target = (output or skill.parent / f"{skill.name}-deep").resolve()
|
|
if target.exists() and any(target.iterdir()):
|
|
raise ValueError(f"deep run directory is not empty: {target}")
|
|
target.mkdir(parents=True, exist_ok=True)
|
|
shutil.copytree(skill, target / "S_fast")
|
|
(target / ".tmp").mkdir()
|
|
(target / "levels").mkdir()
|
|
shutil.copytree(target / "S_fast", target / ".tmp" / "current")
|
|
atomic_write_jsonl(target / "input-traces.jsonl", [asdict(trace) for trace in context.traces])
|
|
state = {
|
|
"run_id": uuid.uuid4().hex[:12],
|
|
"status": "created",
|
|
"model": model,
|
|
"skill_name": skill.name,
|
|
"task": {
|
|
"name": context.task_name,
|
|
"directory": str(context.task_dir),
|
|
"agent": context.agent,
|
|
"model": context.model,
|
|
"prompt": context.prompt,
|
|
},
|
|
"current_rollouts": str(traces.resolve()),
|
|
"current_rollout_skill_sha256": sha256_file(skill / "SKILL.md"),
|
|
"levels": {},
|
|
}
|
|
atomic_write_json(target / "run.json", state)
|
|
return cls(target, analyzer=analyzer, runner=runner)
|
|
|
|
def _save(self) -> None:
|
|
atomic_write_json(self.state_path, self.state)
|
|
|
|
def _prompt(self) -> str:
|
|
return str(self.state["task"]["prompt"])
|
|
|
|
def _current_text(self) -> str:
|
|
return (self.current / "SKILL.md").read_text(encoding="utf-8")
|
|
|
|
def _load_or_rollout(
|
|
self,
|
|
package: Path,
|
|
trace_path: Path,
|
|
batch_id: str,
|
|
seed: list[RolloutTrace] | None = None,
|
|
) -> list[RolloutTrace]:
|
|
if trace_path.is_file():
|
|
return [_trace_from_dict(item) for item in read_jsonl(trace_path)]
|
|
if seed is not None:
|
|
traces = seed
|
|
else:
|
|
_log(f"Starting {ROLLOUTS} rollouts for {batch_id}")
|
|
traces = self.runner.run_batch(
|
|
package,
|
|
self._prompt(),
|
|
batch_id,
|
|
str(self.state["task"]["name"]),
|
|
ROLLOUTS,
|
|
progress=lambda done, total, trace_id: _log(
|
|
f"Rollout {done}/{total} complete: {trace_id}"
|
|
),
|
|
)
|
|
atomic_write_jsonl(trace_path, [asdict(trace) for trace in traces])
|
|
return traces
|
|
|
|
def _load_or_matrix(
|
|
self,
|
|
path: Path,
|
|
level: str,
|
|
units: list[SkillUnit],
|
|
traces: list[RolloutTrace],
|
|
) -> ScoreMatrix:
|
|
value = load_json(path)
|
|
if isinstance(value, dict):
|
|
return ScoreMatrix.from_dict(value)
|
|
columns_dir = path.parent / "matrix-columns"
|
|
cached_columns: dict[str, dict[str, CellScore]] = {}
|
|
for dimension in DIMENSIONS:
|
|
cached = load_json(columns_dir / f"{dimension.lower().replace(' ', '-')}.json")
|
|
if isinstance(cached, dict):
|
|
cached_columns[dimension] = _column_from_dict(cached, traces)
|
|
|
|
def save_column(dimension: str, column: dict[str, CellScore]) -> None:
|
|
atomic_write_json(
|
|
columns_dir / f"{dimension.lower().replace(' ', '-')}.json",
|
|
_column_to_dict(column),
|
|
)
|
|
|
|
_log(f"Scoring full {level} matrix with {self.model}")
|
|
matrix = self.analyzer.score_matrix(
|
|
self._current_text(), self._prompt(), units, traces, level,
|
|
existing_columns=cached_columns,
|
|
result_callback=save_column,
|
|
)
|
|
atomic_write_json(path, matrix.to_dict())
|
|
return matrix
|
|
|
|
def _refresh_matrix(
|
|
self,
|
|
level: str,
|
|
units: list[SkillUnit],
|
|
traces: list[RolloutTrace],
|
|
) -> ScoreMatrix:
|
|
level_state = self.state["levels"][level]
|
|
number = int(level_state.get("refreshes", 0)) + 1
|
|
path = (
|
|
self.run_dir
|
|
/ "levels"
|
|
/ level
|
|
/ "refreshes"
|
|
/ f"refresh-{number:02d}"
|
|
/ "matrix.json"
|
|
)
|
|
_log(f"Refreshing full {level} matrix")
|
|
matrix = self._load_or_matrix(path, level, units, traces)
|
|
atomic_write_json(self.run_dir / "levels" / level / "matrix.json", matrix.to_dict())
|
|
level_state["refreshes"] = number
|
|
level_state.pop("active_dimension", None)
|
|
self._save()
|
|
return matrix
|
|
|
|
def _decisions(self, level: str | None = None) -> list[dict[str, Any]]:
|
|
roots = (
|
|
[self.run_dir / "levels" / level]
|
|
if level else list((self.run_dir / "levels").glob("*"))
|
|
)
|
|
decisions = []
|
|
for root in roots:
|
|
for path in sorted((root / "iterations").glob("iteration-*/decision.json")):
|
|
value = load_json(path)
|
|
if isinstance(value, dict):
|
|
decisions.append(value)
|
|
return decisions
|
|
|
|
def _rejected(self, coordinate: Coordinate) -> list[dict[str, Any]]:
|
|
return [
|
|
decision["rejected_edit"]
|
|
for decision in self._decisions()
|
|
if not decision["accepted"]
|
|
and decision["rejected_edit"]["unit_id"] == coordinate.unit_id
|
|
and decision["rejected_edit"]["dimension"] == coordinate.dimension
|
|
]
|
|
|
|
def _exhausted(self, level: str) -> set[tuple[str, str]]:
|
|
counts: dict[tuple[str, str], int] = {}
|
|
for decision in self._decisions(level):
|
|
if decision["accepted"]:
|
|
continue
|
|
rejected = decision["rejected_edit"]
|
|
key = (rejected["unit_id"], rejected["dimension"])
|
|
counts[key] = counts.get(key, 0) + 1
|
|
return {key for key, count in counts.items() if count >= REJECTION_LIMIT}
|
|
|
|
@staticmethod
|
|
def _block_dimension(level_state: dict[str, Any], dimension: str) -> None:
|
|
blocked = set(level_state.get("blocked_dimensions", []))
|
|
blocked.add(dimension)
|
|
level_state["blocked_dimensions"] = sorted(blocked)
|
|
|
|
@staticmethod
|
|
def _candidate_units(
|
|
units: list[SkillUnit], target: SkillUnit, new_text: str
|
|
) -> list[SkillUnit]:
|
|
delta = len(new_text) - len(target.text)
|
|
updated = []
|
|
for unit in units:
|
|
value = replace(unit)
|
|
if unit.unit_id == target.unit_id:
|
|
value.text = new_text
|
|
value.end = value.start + len(new_text)
|
|
elif unit.start >= target.end:
|
|
value.start += delta
|
|
value.end += delta
|
|
updated.append(value)
|
|
return updated
|
|
|
|
def _apply_edit(
|
|
self,
|
|
unit: SkillUnit,
|
|
edit: LocalEdit,
|
|
candidate: Path,
|
|
section_depth: int | None,
|
|
) -> None:
|
|
self._validate_edit_candidate(unit, edit, section_depth)
|
|
text = self._current_text()
|
|
changed = replace_unit_text(text, unit, edit.new_text)
|
|
temporary = candidate.with_name(f".{candidate.name}.{uuid.uuid4().hex}.tmp")
|
|
shutil.copytree(self.current, temporary)
|
|
atomic_write_text(temporary / "SKILL.md", changed)
|
|
if candidate.exists():
|
|
shutil.rmtree(candidate)
|
|
candidate.parent.mkdir(parents=True, exist_ok=True)
|
|
temporary.rename(candidate)
|
|
|
|
def _validate_edit_candidate(
|
|
self,
|
|
unit: SkillUnit,
|
|
edit: LocalEdit,
|
|
section_depth: int | None,
|
|
) -> None:
|
|
"""Validate a local edit against both unit and whole-document invariants."""
|
|
|
|
validate_edit(unit, edit.new_text, section_depth)
|
|
text = self._current_text()
|
|
if text[unit.start:unit.end] != unit.text:
|
|
raise ValueError("target unit no longer matches current SKILL.md")
|
|
changed = replace_unit_text(text, unit, edit.new_text)
|
|
before = [(item.heading, item.heading_depth) for item in parse_sections(text)]
|
|
after = [(item.heading, item.heading_depth) for item in parse_sections(changed)]
|
|
if before != after:
|
|
raise ValueError("local edit changed section boundaries")
|
|
|
|
@staticmethod
|
|
def _validation_feedback(
|
|
attempts: list[dict[str, Any]],
|
|
) -> list[dict[str, Any]]:
|
|
return [
|
|
{
|
|
"edit_summary": str(item.get("edit_summary", "invalid generated edit")),
|
|
"reject_reason": f"structural_validation_failed: {item['error']}",
|
|
"new_text_hash": str(item.get("new_text_hash", "")),
|
|
}
|
|
for item in attempts
|
|
]
|
|
|
|
def _valid_edit_or_rejection(
|
|
self,
|
|
*,
|
|
level: str,
|
|
number: int,
|
|
iteration_dir: Path,
|
|
unit: SkillUnit,
|
|
coordinate: Coordinate,
|
|
cell: CellScore,
|
|
section_depth: int | None,
|
|
) -> tuple[LocalEdit | None, bool, list[dict[str, Any]]]:
|
|
"""Load or generate a valid edit, feeding structural failures back to the model."""
|
|
|
|
edit_path = iteration_dir / "edit.json"
|
|
attempts_path = iteration_dir / "edit-attempts.json"
|
|
attempts_value = load_json(attempts_path, [])
|
|
attempts = attempts_value if isinstance(attempts_value, list) else []
|
|
cached_value = load_json(edit_path)
|
|
|
|
if isinstance(cached_value, dict):
|
|
try:
|
|
cached = LocalEdit.from_dict(cached_value)
|
|
cached.new_text = preserve_unit_boundary(unit, cached.new_text)
|
|
self._validate_edit_candidate(unit, cached, section_depth)
|
|
return cached, False, attempts
|
|
except ValueError as exc:
|
|
text = str(cached_value.get("new_text", ""))
|
|
attempts.append({
|
|
"source": "cached",
|
|
"error": str(exc),
|
|
"edit_summary": str(cached_value.get("edit_summary", "")),
|
|
"new_text_hash": sha256_text(text) if text else "",
|
|
})
|
|
atomic_write_json(attempts_path, attempts)
|
|
_log(
|
|
f"{level} iteration {number}: cached local edit is invalid: {exc}; "
|
|
"regenerating"
|
|
)
|
|
|
|
for attempt in range(1, MAX_EDIT_GENERATION_ATTEMPTS + 1):
|
|
_log(
|
|
f"{level} iteration {number}: generating local edit for "
|
|
f"{coordinate.unit_id}/{coordinate.dimension} "
|
|
f"(attempt {attempt}/{MAX_EDIT_GENERATION_ATTEMPTS})"
|
|
)
|
|
feedback = self._rejected(coordinate) + self._validation_feedback(attempts)
|
|
edit: LocalEdit | None = None
|
|
try:
|
|
edit = self.analyzer.generate_edit(
|
|
coordinate,
|
|
unit,
|
|
cell,
|
|
feedback,
|
|
self._prompt(),
|
|
)
|
|
edit.new_text = preserve_unit_boundary(unit, edit.new_text)
|
|
self._validate_edit_candidate(unit, edit, section_depth)
|
|
except ValueError as exc:
|
|
text = edit.new_text if edit is not None else ""
|
|
attempts.append({
|
|
"source": "generated",
|
|
"generation_attempt": attempt,
|
|
"error": str(exc),
|
|
"edit_summary": edit.edit_summary if edit is not None else "",
|
|
"new_text_hash": sha256_text(text) if text else "",
|
|
})
|
|
atomic_write_json(attempts_path, attempts)
|
|
_log(
|
|
f"{level} iteration {number}: local edit validation failed "
|
|
f"(attempt {attempt}/{MAX_EDIT_GENERATION_ATTEMPTS}): {exc}"
|
|
)
|
|
continue
|
|
|
|
assert edit is not None
|
|
atomic_write_json(edit_path, asdict(edit))
|
|
return edit, True, attempts
|
|
|
|
return None, True, attempts
|
|
|
|
def _commit_iteration(
|
|
self,
|
|
level: str,
|
|
number: int,
|
|
iteration_dir: Path,
|
|
matrix: ScoreMatrix,
|
|
current_trace_path: Path,
|
|
) -> ScoreMatrix:
|
|
decision = load_json(iteration_dir / "decision.json")
|
|
if not isinstance(decision, dict):
|
|
raise ValueError("missing iteration decision")
|
|
if decision["accepted"]:
|
|
candidate = iteration_dir / "candidate" / self.state["skill_name"]
|
|
replacement = self.temp / "next-current"
|
|
shutil.rmtree(replacement, ignore_errors=True)
|
|
shutil.copytree(candidate, replacement)
|
|
shutil.rmtree(self.current)
|
|
replacement.rename(self.current)
|
|
matrix = ScoreMatrix.from_dict(decision["matrix_after"])
|
|
atomic_write_json(self.run_dir / "levels" / level / "matrix.json", matrix.to_dict())
|
|
candidate_traces = read_jsonl(iteration_dir / "candidate-traces.jsonl")
|
|
atomic_write_jsonl(current_trace_path, candidate_traces)
|
|
rollout_output = decision.get("rollout_output")
|
|
rollout_skill_sha256 = decision.get("rollout_skill_sha256")
|
|
if isinstance(rollout_output, str) and isinstance(rollout_skill_sha256, str):
|
|
self.state["current_rollouts"] = rollout_output
|
|
self.state["current_rollout_skill_sha256"] = rollout_skill_sha256
|
|
else:
|
|
self.state.pop("current_rollouts", None)
|
|
self.state.pop("current_rollout_skill_sha256", None)
|
|
level_state = self.state["levels"][level]
|
|
if int(level_state.get("iterations", 0)) < number:
|
|
level_state["iterations"] = number
|
|
self._save()
|
|
return matrix
|
|
|
|
def _run_iteration(
|
|
self,
|
|
level: str,
|
|
number: int,
|
|
level_dir: Path,
|
|
matrix: ScoreMatrix,
|
|
coordinate: Coordinate,
|
|
current_trace_path: Path,
|
|
section_depth: int | None,
|
|
) -> ScoreMatrix:
|
|
iteration_dir = level_dir / "iterations" / f"iteration-{number:02d}"
|
|
iteration_dir.mkdir(parents=True, exist_ok=True)
|
|
unit = matrix.unit(coordinate.unit_id)
|
|
edit, regenerated, validation_attempts = self._valid_edit_or_rejection(
|
|
level=level,
|
|
number=number,
|
|
iteration_dir=iteration_dir,
|
|
unit=unit,
|
|
coordinate=coordinate,
|
|
cell=matrix.columns[coordinate.dimension][coordinate.unit_id],
|
|
section_depth=section_depth,
|
|
)
|
|
if edit is None:
|
|
decision = {
|
|
"coordinate": asdict(coordinate),
|
|
"accepted": False,
|
|
"reason": "edit_validation_exhausted",
|
|
"target_delta": 0.0,
|
|
"validation_attempts": validation_attempts,
|
|
"rejected_edit": {
|
|
"unit_id": coordinate.unit_id,
|
|
"dimension": coordinate.dimension,
|
|
"edit_summary": (
|
|
"Could not generate a structurally valid local edit after "
|
|
f"{MAX_EDIT_GENERATION_ATTEMPTS} attempts."
|
|
),
|
|
"score_change": 0.0,
|
|
"reject_reason": "edit_validation_exhausted",
|
|
"new_text_hash": str(
|
|
validation_attempts[-1].get("new_text_hash", "")
|
|
) if validation_attempts else "",
|
|
},
|
|
}
|
|
atomic_write_json(iteration_dir / "decision.json", decision)
|
|
_log(
|
|
f"{level} iteration {number}: local edit validation exhausted; "
|
|
"recording rejection and continuing"
|
|
)
|
|
return self._commit_iteration(
|
|
level, number, iteration_dir, matrix, current_trace_path
|
|
)
|
|
edit_hash = sha256_text(edit.new_text)
|
|
if any(item["new_text_hash"] == edit_hash for item in self._rejected(coordinate)):
|
|
decision = {
|
|
"coordinate": asdict(coordinate),
|
|
"accepted": False,
|
|
"reason": "exact_duplicate_rejected_edit",
|
|
"target_delta": 0.0,
|
|
"rejected_edit": {
|
|
"unit_id": coordinate.unit_id,
|
|
"dimension": coordinate.dimension,
|
|
"edit_summary": edit.edit_summary,
|
|
"score_change": 0.0,
|
|
"reject_reason": "exact_duplicate_rejected_edit",
|
|
"new_text_hash": edit_hash,
|
|
},
|
|
}
|
|
atomic_write_json(iteration_dir / "decision.json", decision)
|
|
return self._commit_iteration(level, number, iteration_dir, matrix, current_trace_path)
|
|
candidate = iteration_dir / "candidate" / self.state["skill_name"]
|
|
if regenerated:
|
|
shutil.rmtree(candidate, ignore_errors=True)
|
|
for stale in (
|
|
iteration_dir / "candidate-traces.jsonl",
|
|
iteration_dir / "comparison.json",
|
|
):
|
|
if stale.exists():
|
|
stale.unlink()
|
|
if not (candidate / "SKILL.md").is_file():
|
|
self._apply_edit(unit, edit, candidate, section_depth)
|
|
candidate_units = self._candidate_units(matrix.units, unit, edit.new_text)
|
|
candidate_skill_sha256 = sha256_file(candidate / "SKILL.md")
|
|
batch_id = (
|
|
f"{self.state['run_id']}-deep-{level}-i{number:02d}-"
|
|
f"{candidate_skill_sha256[:12]}"
|
|
)
|
|
traces = self._load_or_rollout(
|
|
candidate,
|
|
iteration_dir / "candidate-traces.jsonl",
|
|
batch_id,
|
|
)
|
|
comparison_path = iteration_dir / "comparison.json"
|
|
comparison = load_json(comparison_path)
|
|
if not isinstance(comparison, dict):
|
|
_log(
|
|
f"{level} iteration {number}: comparing "
|
|
f"{coordinate.unit_id}/{coordinate.dimension}"
|
|
)
|
|
incumbent_traces = [
|
|
_trace_from_dict(item) for item in read_jsonl(current_trace_path)
|
|
]
|
|
comparison = self.analyzer.compare_cell(
|
|
self._prompt(),
|
|
unit,
|
|
next(item for item in candidate_units if item.unit_id == unit.unit_id),
|
|
incumbent_traces,
|
|
traces,
|
|
coordinate.dimension,
|
|
)
|
|
incumbent_timeout_rate = (
|
|
sum(trace.timed_out is True for trace in incumbent_traces)
|
|
/ len(incumbent_traces)
|
|
)
|
|
candidate_timeout_rate = (
|
|
sum(trace.timed_out is True for trace in traces) / len(traces)
|
|
)
|
|
if (
|
|
candidate_timeout_rate > incumbent_timeout_rate
|
|
and "timeout" not in comparison["runtime_regressions"]
|
|
):
|
|
comparison["runtime_regressions"].append("timeout")
|
|
atomic_write_json(comparison_path, comparison)
|
|
|
|
delta = float(comparison["candidate_score"]) - float(
|
|
comparison["incumbent_score"]
|
|
)
|
|
if not comparison["task_relevant"]:
|
|
accepted, reason = False, "target_unit_not_task_relevant"
|
|
elif comparison["runtime_regressions"]:
|
|
accepted = False
|
|
reason = "runtime_regressed:" + ",".join(comparison["runtime_regressions"])
|
|
elif delta < 0.5:
|
|
accepted, reason = False, "target_cell_did_not_improve"
|
|
else:
|
|
accepted, reason = True, "target_improved_without_runtime_regression"
|
|
decision: dict[str, Any] = {
|
|
"coordinate": asdict(coordinate),
|
|
"accepted": accepted,
|
|
"reason": reason,
|
|
"target_delta": delta,
|
|
"rollout_skill_sha256": candidate_skill_sha256,
|
|
}
|
|
artifacts_dir = getattr(self.runner, "artifacts_dir", None)
|
|
rollout_output = artifacts_dir(batch_id) if callable(artifacts_dir) else None
|
|
if isinstance(rollout_output, Path):
|
|
decision["rollout_output"] = str(rollout_output.resolve())
|
|
if accepted:
|
|
updated = ScoreMatrix(matrix.level, candidate_units, dict(matrix.columns))
|
|
updated.columns[coordinate.dimension] = dict(
|
|
matrix.columns[coordinate.dimension]
|
|
)
|
|
updated.columns[coordinate.dimension][coordinate.unit_id] = CellScore(
|
|
float(comparison["candidate_score"]),
|
|
list(comparison["candidate_evidence"]),
|
|
str(comparison["reason"]),
|
|
)
|
|
decision["matrix_after"] = updated.to_dict()
|
|
else:
|
|
decision["rejected_edit"] = {
|
|
"unit_id": coordinate.unit_id,
|
|
"dimension": coordinate.dimension,
|
|
"edit_summary": edit.edit_summary,
|
|
"score_change": delta,
|
|
"reject_reason": reason,
|
|
"new_text_hash": edit_hash,
|
|
}
|
|
atomic_write_json(iteration_dir / "decision.json", decision)
|
|
return self._commit_iteration(level, number, iteration_dir, matrix, current_trace_path)
|
|
|
|
def _run_level(
|
|
self,
|
|
level: str,
|
|
units: list[SkillUnit],
|
|
seed_traces: list[RolloutTrace] | None = None,
|
|
section_depth: int | None = None,
|
|
) -> tuple[ScoreMatrix, list[RolloutTrace], Coordinate | None]:
|
|
level_dir = self.run_dir / "levels" / level
|
|
level_dir.mkdir(parents=True, exist_ok=True)
|
|
level_state = self.state["levels"].setdefault(level, {
|
|
"iterations": 0, "completed": False,
|
|
})
|
|
current_trace_path = level_dir / "current-traces.jsonl"
|
|
traces = self._load_or_rollout(
|
|
self.current,
|
|
current_trace_path,
|
|
f"{self.state['run_id']}-deep-{level}-initial",
|
|
seed=seed_traces,
|
|
)
|
|
matrix = self._load_or_matrix(level_dir / "matrix.json", level, units, traces)
|
|
if level_state.get("completed"):
|
|
exhausted_value = level_state.get("exhausted_coordinate")
|
|
exhausted = Coordinate(**exhausted_value) if isinstance(exhausted_value, dict) else None
|
|
return matrix, traces, exhausted
|
|
exhausted_coordinate: Coordinate | None = None
|
|
max_iterations = (
|
|
MAX_SECTION_ITERATIONS if level == "section" else MAX_PARAGRAPH_ITERATIONS
|
|
)
|
|
while int(level_state["iterations"]) < max_iterations:
|
|
excluded = self._exhausted(level)
|
|
excluded.update(
|
|
(unit.unit_id, dimension)
|
|
for dimension in level_state.get("blocked_dimensions", [])
|
|
for unit in matrix.units
|
|
)
|
|
pending_number = int(level_state["iterations"]) + 1
|
|
pending_dir = level_dir / "iterations" / f"iteration-{pending_number:02d}"
|
|
pending_decision = load_json(pending_dir / "decision.json")
|
|
if isinstance(pending_decision, dict):
|
|
pending_coordinate = Coordinate(**pending_decision["coordinate"])
|
|
level_state.setdefault("active_dimension", pending_coordinate.dimension)
|
|
matrix = self._commit_iteration(
|
|
level, pending_number, pending_dir, matrix, current_trace_path
|
|
)
|
|
traces = [_trace_from_dict(item) for item in read_jsonl(current_trace_path)]
|
|
if len(self._rejected(pending_coordinate)) >= REJECTION_LIMIT:
|
|
exhausted_coordinate = pending_coordinate
|
|
level_state["stop_reason"] = "coordinate_exhausted"
|
|
level_state["exhausted_coordinate"] = asdict(pending_coordinate)
|
|
break
|
|
if matrix.select_coordinate(
|
|
GAP_THRESHOLD, level_state.get("active_dimension"), excluded
|
|
) is None:
|
|
matrix = self._refresh_matrix(level, matrix.units, traces)
|
|
continue
|
|
active_dimension = level_state.get("active_dimension")
|
|
coordinate = matrix.select_coordinate(
|
|
GAP_THRESHOLD, active_dimension, excluded
|
|
)
|
|
if coordinate is None:
|
|
if active_dimension is None:
|
|
level_state["stop_reason"] = (
|
|
"normalized_gap_converged"
|
|
if max(matrix.normalized_gaps().values(), default=0.0) <= GAP_THRESHOLD
|
|
else "available_coordinates_exhausted"
|
|
)
|
|
break
|
|
matrix = self._refresh_matrix(level, matrix.units, traces)
|
|
continue
|
|
if active_dimension is None:
|
|
level_state["active_dimension"] = coordinate.dimension
|
|
self._save()
|
|
number = int(level_state["iterations"]) + 1
|
|
matrix = self._run_iteration(
|
|
level, number, level_dir, matrix, coordinate,
|
|
current_trace_path, section_depth,
|
|
)
|
|
traces = [_trace_from_dict(item) for item in read_jsonl(current_trace_path)]
|
|
if len(self._rejected(coordinate)) >= REJECTION_LIMIT:
|
|
exhausted_coordinate = coordinate
|
|
level_state["stop_reason"] = "coordinate_exhausted"
|
|
level_state["exhausted_coordinate"] = asdict(coordinate)
|
|
break
|
|
else:
|
|
decisions = self._decisions(level)
|
|
last = decisions[-1] if decisions else {}
|
|
if matrix.select_coordinate(GAP_THRESHOLD) is None:
|
|
level_state["stop_reason"] = "normalized_gap_converged"
|
|
else:
|
|
level_state["stop_reason"] = (
|
|
"max_iterations_after_accept" if last.get("accepted") else "max_iterations"
|
|
)
|
|
level_state["completed"] = True
|
|
self._save()
|
|
return matrix, traces, exhausted_coordinate
|
|
|
|
def drive(self) -> Path:
|
|
if self.state.get("status") == "complete" and (self.run_dir / "S_final").is_dir():
|
|
return self.run_dir / "S_final"
|
|
_log(f"Deep Loop start/resume: {self.run_dir}")
|
|
input_traces = [
|
|
_trace_from_dict(item) for item in read_jsonl(self.run_dir / "input-traces.jsonl")
|
|
]
|
|
while True:
|
|
section_units = parse_sections(self._current_text())
|
|
section_matrix, traces, exhausted = self._run_level(
|
|
"section", section_units, seed_traces=input_traces
|
|
)
|
|
if exhausted is None:
|
|
break
|
|
target_score = section_matrix.columns[exhausted.dimension][exhausted.unit_id].score
|
|
current_sections = parse_sections(self._current_text())
|
|
section = next((item for item in current_sections if item.unit_id == exhausted.unit_id), None)
|
|
paragraphs = parse_paragraphs(section) if section is not None else []
|
|
section_state = self.state["levels"]["section"]
|
|
if (
|
|
target_score > 3.5
|
|
or len(paragraphs) < 2
|
|
or section_state.get("paragraph_returned")
|
|
):
|
|
self._block_dimension(section_state, exhausted.dimension)
|
|
section_state["completed"] = False
|
|
for key in ("stop_reason", "exhausted_coordinate", "active_dimension"):
|
|
section_state.pop(key, None)
|
|
self._save()
|
|
continue
|
|
_log(f"Descending into paragraphs of {exhausted.unit_id}")
|
|
self.state["levels"].setdefault(
|
|
"paragraph", {"iterations": 0, "completed": False}
|
|
)["active_dimension"] = exhausted.dimension
|
|
self._save()
|
|
_, paragraph_traces, _ = self._run_level(
|
|
"paragraph", paragraphs, seed_traces=traces,
|
|
section_depth=section.heading_depth if section else None,
|
|
)
|
|
atomic_write_jsonl(
|
|
self.run_dir / "levels" / "section" / "current-traces.jsonl",
|
|
[asdict(trace) for trace in paragraph_traces],
|
|
)
|
|
section_state["completed"] = False
|
|
section_state["paragraph_returned"] = True
|
|
self._block_dimension(section_state, exhausted.dimension)
|
|
for key in ("stop_reason", "exhausted_coordinate", "active_dimension"):
|
|
section_state.pop(key, None)
|
|
self._refresh_matrix(
|
|
"section", parse_sections(self._current_text()), paragraph_traces
|
|
)
|
|
return self._complete()
|
|
|
|
def _complete(self) -> Path:
|
|
target = self.run_dir / "S_final"
|
|
if target.exists():
|
|
shutil.rmtree(target)
|
|
shutil.copytree(self.current, target)
|
|
final_skill_hash = sha256_file(target / "SKILL.md")
|
|
final_rollouts = self.state.get("current_rollouts")
|
|
final_rollout_hash = self.state.get("current_rollout_skill_sha256")
|
|
if (
|
|
not isinstance(final_rollouts, str)
|
|
or not Path(final_rollouts).is_dir()
|
|
or final_rollout_hash != final_skill_hash
|
|
):
|
|
final_rollouts = None
|
|
report = {
|
|
"status": "complete",
|
|
"input_package_hash": package_hash(self.run_dir / "S_fast"),
|
|
"final_package_hash": package_hash(target),
|
|
"final_rollouts": final_rollouts,
|
|
"final_rollout_skill_sha256": final_skill_hash if final_rollouts else None,
|
|
"task": self.state["task"]["name"],
|
|
"agent": self.state["task"]["agent"],
|
|
"target_model": self.state["task"]["model"],
|
|
"model": self.model,
|
|
"rollout_backend": "skillsbench_development",
|
|
"production_rollout_backend": "blank_container_required",
|
|
"verifier_signal_used": False,
|
|
"levels": self.state["levels"],
|
|
"decisions": {
|
|
name: self._decisions(name)
|
|
for name in ("section", "paragraph")
|
|
if name in self.state["levels"]
|
|
},
|
|
}
|
|
atomic_write_json(self.run_dir / "report.json", report)
|
|
self.state["status"] = "complete"
|
|
self._save()
|
|
shutil.rmtree(self.temp, ignore_errors=True)
|
|
_log(f"Deep Loop complete: {target}")
|
|
return target
|