Files
2026-09-04 14:58:42 +08:00

320 lines
11 KiB
Python

"""动态编译主流水线;本模块只负责阶段编排。"""
from __future__ import annotations
import os
import sys
from pathlib import Path
from typing import Any
from .cache import (
SCORE_ARTIFACTS,
fingerprint,
load_candidate_cache,
load_maps_cache,
load_reduction_cache,
load_score_cache,
trace_fingerprint,
write_manifest,
)
from .optimization.analyzer import SemanticAnalyzer, SemanticClient
from .optimization.patch import apply_patches
from .optimization.selection import relative_high_low
from .paths import default_outputs, project_path
from .scoring.agentrm import (
DEFAULT_BATCH_SIZE,
DEFAULT_CONCURRENCY,
DEFAULT_MAX_LENGTH,
DEFAULT_RM_API_URL,
DEFAULT_TIMEOUT,
AgentRM,
)
from .scoring.service import TraceScorer
from .scoring.pre_score import PRE_SCORE_VERSION, PreScorer, RelevanceJudge
from .storage import (
atomic_write_json,
package_hash,
read_jsonl,
sha256_file,
)
from .traces.benchflow import load_benchflow_traces
GROUP_SIZE = 3
SCORE_VERSION = 1
MAP_VERSION = 1
REDUCTION_VERSION = 1
PATCH_VERSION = 1
def _log(message: str) -> None:
print(f"[dynamic_compile.fast] {message}", file=sys.stderr, flush=True)
def _implementation_name(value: object | None, default: str) -> str:
if value is None:
return default
return type(value).__module__ + "." + type(value).__qualname__
def run_pipeline(
trace_input: Path,
skill_package: Path,
*,
score_output: Path | None = None,
output: Path | None = None,
model: str = "opencode/deepseek-v4-pro",
max_parallel: int = 3,
rm_api_url: str | None = None,
rm_max_length: int = DEFAULT_MAX_LENGTH,
rm_timeout: float = DEFAULT_TIMEOUT,
rm_concurrency: int = DEFAULT_CONCURRENCY,
rm_batch_size: int = DEFAULT_BATCH_SIZE,
force: bool = False,
analyzer: SemanticAnalyzer | None = None,
pre_scorer: PreScorer | None = None,
agentrm: AgentRM | None = None,
) -> Path:
"""从历史 BenchFlow 轨迹生成一个候选 Skill 包。"""
trace_input = project_path(trace_input)
skill_package = project_path(skill_package)
if not (skill_package / "SKILL.md").is_file():
raise ValueError("skill package must contain a root SKILL.md")
if max_parallel < 1:
raise ValueError("max_parallel must be at least 1")
if agentrm is None and min(
rm_max_length, rm_timeout, rm_concurrency, rm_batch_size
) <= 0:
raise ValueError("AgentRM numeric options must be positive")
traces = load_benchflow_traces(trace_input)
identities = {(trace.task_name, trace.compile_type) for trace in traces}
if len(identities) != 1:
raise ValueError(
f"trace input must contain one task and compile type: {sorted(identities)}"
)
if len(traces) < GROUP_SIZE * 2:
raise ValueError(f"at least {GROUP_SIZE * 2} traces are required")
if len({trace.test_name for trace in traces}) != len(traces):
raise ValueError("trace input contains duplicate test names")
_, compile_type = next(iter(identities))
if score_output is None or output is None:
default_score, default_output = default_outputs(trace_input, compile_type)
score_dir = project_path(score_output) if score_output else default_score
output_dir = project_path(output) if output else default_output
if output_dir == skill_package or output_dir.is_relative_to(skill_package):
raise ValueError("output directory must not be inside the input skill package")
score_dir.mkdir(parents=True, exist_ok=True)
output_dir.mkdir(parents=True, exist_ok=True)
traces_hash = trace_fingerprint(traces)
resolved_rm_url = rm_api_url or os.environ.get("RM_API_URL", DEFAULT_RM_API_URL)
score_config: dict[str, Any] = {
"score_version": SCORE_VERSION,
"pre_score_version": PRE_SCORE_VERSION,
"model": model,
"pre_scorer": _implementation_name(pre_scorer, "PreScorer/RelevanceJudge"),
"agentrm": _implementation_name(agentrm, "AgentRM/HttpAgentRMBackend"),
"rm_api_url": resolved_rm_url,
"rm_max_length": rm_max_length,
"rm_timeout": rm_timeout,
"rm_concurrency": rm_concurrency,
"rm_batch_size": rm_batch_size,
}
score_input_hash = fingerprint({"traces": traces_hash, "config": score_config})
_log(f"loaded {len(traces)} traces from {trace_input}")
cached_score = load_score_cache(
traces, score_dir, input_hash=score_input_hash, config=score_config
)
if cached_score is not None:
scores, score_rows = cached_score
_log("reusing complete, input-matched scoring cache")
else:
scoring = TraceScorer(
pre_scorer or PreScorer(RelevanceJudge(model=model), max_parallel),
agentrm
or AgentRM(
api_url=resolved_rm_url,
max_length=rm_max_length,
timeout=rm_timeout,
concurrency=rm_concurrency,
batch_size=rm_batch_size,
),
)
_log("scoring traces")
scores = scoring.score_all(traces, score_dir)
score_rows = read_jsonl(score_dir / "effective_scores.jsonl")
write_manifest(
score_dir,
".score-cache.json",
stage="score",
input_hash=score_input_hash,
config=score_config,
artifacts=SCORE_ARTIFACTS,
)
high, low = relative_high_low(scores, count=GROUP_SIZE)
selected = high + low
semantic: SemanticAnalyzer | None = analyzer
def get_semantic() -> SemanticAnalyzer:
nonlocal semantic
if semantic is None:
semantic = SemanticAnalyzer(SemanticClient(model), max_parallel)
return semantic
semantic_implementation = _implementation_name(analyzer, "SemanticAnalyzer/SemanticClient")
map_config = {
"map_version": MAP_VERSION,
"model": model,
"semantic_implementation": semantic_implementation,
"max_parallel": max_parallel,
}
map_input_hash = fingerprint(
{
"traces": traces_hash,
"selected": selected,
"scores": {trace_id: scores[trace_id] for trace_id in selected},
}
)
maps = load_maps_cache(
output_dir,
selected,
input_hash=map_input_hash,
config=map_config,
)
if maps is not None:
_log(f"reusing complete Top {GROUP_SIZE} / Bottom {GROUP_SIZE} Map cache")
else:
_log(f"mapping Top {GROUP_SIZE} / Bottom {GROUP_SIZE} traces")
selected_set = set(selected)
maps = get_semantic().map_all(
[trace for trace in traces if trace.trace_id in selected_set],
scores,
set(high),
set(low),
progress=lambda done, total, trace_id: _log(
f"Map {done}/{total}: {trace_id}"
),
)
atomic_write_json(output_dir / "maps.json", maps)
write_manifest(
output_dir,
".maps-cache.json",
stage="maps",
input_hash=map_input_hash,
config=map_config,
artifacts=("maps.json",),
)
maps_by_id = {str(item["trace_id"]): item for item in maps}
scores_by_id = {str(item["trace_id"]): item for item in score_rows}
skill_path = skill_package / "SKILL.md"
skill_text = skill_path.read_text(encoding="utf-8")
skill_hash = sha256_file(skill_path)
skill_package_hash = package_hash(skill_package)
reduction_config = {
"reduction_version": REDUCTION_VERSION,
"model": model,
"semantic_implementation": semantic_implementation,
}
reduction_input_hash = fingerprint(
{
"traces": traces_hash,
"package": skill_package_hash,
"maps": [maps_by_id[trace_id] for trace_id in selected],
"score_rows": [scores_by_id[trace_id] for trace_id in selected],
"high": high,
"low": low,
}
)
reduction = None if force else load_reduction_cache(
output_dir,
input_hash=reduction_input_hash,
config=reduction_config,
)
if reduction is not None:
_log("reusing complete Top/Bottom reduction cache")
else:
_log("reducing Top/Bottom contrast")
reduction = get_semantic().reduce(
skill_text,
[maps_by_id[trace_id] for trace_id in selected],
[scores_by_id[trace_id] for trace_id in selected],
high,
low,
[],
)
atomic_write_json(output_dir / "reduction.json", reduction)
write_manifest(
output_dir,
".reduction-cache.json",
stage="reduction",
input_hash=reduction_input_hash,
config=reduction_config,
artifacts=("reduction.json",),
)
candidate = output_dir / "candidate-skill" / skill_package.name
candidate_config = {
"patch_version": PATCH_VERSION,
"model": model,
"semantic_implementation": semantic_implementation,
"attempts": 3,
}
candidate_input_hash = fingerprint(
{
"traces": traces_hash,
"package": skill_package_hash,
"reduction": reduction,
}
)
cached_candidate = None if force else load_candidate_cache(
output_dir,
skill_package.name,
skill_hash,
input_hash=candidate_input_hash,
config=candidate_config,
)
if cached_candidate is not None:
_log(f"reusing complete candidate skill: {cached_candidate}")
return cached_candidate
error = ""
for attempt in range(3):
try:
_log(f"generating patch bundle ({attempt + 1}/3)")
patches = get_semantic().generate_patches(skill_path, reduction, [], error)
candidate_hash = apply_patches(skill_package, candidate, patches)
atomic_write_json(
output_dir / "patch.json",
{
"patches": [patch.to_dict() for patch in patches],
"skill_hash": skill_hash,
"candidate_skill_hash": candidate_hash,
},
)
write_manifest(
output_dir,
".candidate-cache.json",
stage="candidate",
input_hash=candidate_input_hash,
config=candidate_config,
artifacts=(
"patch.json",
f"candidate-skill/{skill_package.name}",
),
)
_log(f"candidate skill ready: {candidate}")
return candidate
except (OSError, ValueError, RuntimeError) as exc:
error = str(exc)
if attempt == 2:
raise RuntimeError(
f"could not generate an applicable patch: {error}"
) from exc
raise AssertionError("unreachable")