| """Phase-separated E02 interactive localization experiment.""" |
|
|
| from __future__ import annotations |
|
|
| from dataclasses import asdict |
| from hashlib import sha256 |
| import json |
| from pathlib import Path |
| import time |
| from typing import Any, Sequence |
|
|
| from .confirmatory_retrieval import retrieve_treatment, subprocess_git |
| from .context_packing import ( |
| EVIDENCE_TOKEN_BUDGET, |
| candidates_from_records, |
| load_ranking_records, |
| pack_role_summaries, |
| pack_skeletons, |
| pack_snippets, |
| pack_specialized_channels, |
| pack_whole_files, |
| ranking_path, |
| ) |
| from .fusion import unique_files |
| from .llm_localization import LocalizationError, parse_selection |
| from .lm_studio import LMStudioClient, LMStudioError |
| from .lm_studio_embeddings import LMStudioEmbeddingClient |
| from .pilot import research_code_revision, retrieval_metrics |
| from .repository import GitSnapshot, chunk_snapshot |
| from .retrieval import BM25FuzzyRetriever, DenseRetriever, ExactRetriever, SQLiteEmbeddingCache |
| from .specs import ( |
| HarnessSpec, |
| TaskSpec, |
| load_embeddings, |
| load_experiments, |
| load_harnesses, |
| load_models, |
| load_task_split, |
| load_tasks, |
| ) |
| from .syntax_index import SymbolGraph, SyntaxRetriever, parse_snapshot |
| from .telemetry import EventWriter, RunIdentity, run_directory |
| from .tokenization import QwenTokenCounter |
| from .vector_backends import FaissFlatRetriever |
|
|
|
|
| LOCALIZE_SYSTEM = """You are performing blinded bug localization in a large Go repository. |
| Use only the issue and evidence supplied by the harness. Select source files likely to |
| require changes. Return exactly one JSON object with keys files and reasoning. files must |
| contain 1-10 distinct evidence paths; reasoning must be below 180 words. Do not write a patch |
| and do not use Markdown.""" |
|
|
| QUERY_SYSTEM = """You are reformulating a code-search query after inspecting initial evidence. |
| Return exactly {"query":"..."} with one concise query that would retrieve missing implementation |
| locations. Do not include reasoning, Markdown, or any other key.""" |
|
|
|
|
| class InteractiveExperimentError(RuntimeError): |
| """Raised when an E02 phase cannot preserve its frozen protocol.""" |
|
|
|
|
| def score_selection_response( |
| response: dict[str, Any], allowed_paths: set[str] |
| ) -> tuple[dict[str, Any], str | None]: |
| """Parse a constrained selection or score an interface violation as no selection. |
| |
| Invalid model output is an experimental outcome, not an infrastructure failure. The |
| unmodified response remains in the run artifacts and the violation is explicitly logged. |
| """ |
|
|
| try: |
| return parse_selection(response, allowed_paths), None |
| except LocalizationError as exc: |
| return {"files": [], "reasoning": ""}, str(exc) |
|
|
|
|
| def exclusive_loaded(native_models: Sequence[dict[str, Any]]) -> tuple[str, ...]: |
| return tuple( |
| str(record.get("key")) for record in native_models if record.get("loaded_instances") |
| ) |
|
|
|
|
| def selected_tasks(root: Path, experiment: Any, task_filter: set[str] | None) -> list[TaskSpec]: |
| catalog = load_tasks(root) |
| split = load_task_split(root / "tasks" / "splits" / f"{experiment.task_split}.txt") |
| return [catalog[item] for item in split if task_filter is None or item in task_filter] |
|
|
|
|
| def initial_context( |
| root: Path, |
| task: TaskSpec, |
| harness: HarnessSpec, |
| snapshot: GitSnapshot, |
| tokenizer: QwenTokenCounter, |
| ) -> tuple[str, tuple[str, ...], int]: |
| if harness.interface == "specialized": |
| return pack_specialized_channels( |
| root, task.task_id, snapshot, task.base_commit, tokenizer |
| ) |
| records = load_ranking_records(ranking_path(root, "E01", "H008", task.task_id)) |
| candidates = candidates_from_records(snapshot, task.base_commit, records) |
| paths = [candidate.path for candidate in candidates] |
| if harness.packing == "ranked_snippets": |
| return pack_snippets(tokenizer, candidates) |
| if harness.packing == "skeletons": |
| return pack_skeletons(tokenizer, snapshot, task.base_commit, paths) |
| if harness.packing == "whole_files": |
| return pack_whole_files(tokenizer, snapshot, task.base_commit, paths) |
| if harness.packing == "role_summaries": |
| return pack_role_summaries(tokenizer, snapshot, task.base_commit, paths) |
| raise InteractiveExperimentError(f"unsupported packing policy {harness.packing}") |
|
|
|
|
| def localization_prompt(task: TaskSpec, context: str, phase: str = "initial") -> str: |
| return f"ISSUE:\n{task.statement}\n\n{phase.upper()} EVIDENCE:\n{context}" |
|
|
|
|
| def parse_query(response: dict[str, Any]) -> str: |
| try: |
| content = response["choices"][0]["message"]["content"] |
| stripped = content.strip() |
| if stripped.startswith("```"): |
| stripped = stripped.removeprefix("```json").removeprefix("```") |
| stripped = stripped.removesuffix("```").strip() |
| try: |
| value = json.loads(stripped) |
| except json.JSONDecodeError: |
| start, end = stripped.find("{"), stripped.rfind("}") |
| value = json.loads(stripped[start : end + 1]) |
| query = value["query"] |
| except (KeyError, IndexError, TypeError, json.JSONDecodeError) as exc: |
| raise InteractiveExperimentError("query-reformulation response is not valid JSON") from exc |
| if set(value) != {"query"} or not isinstance(query, str) or not query.strip(): |
| raise InteractiveExperimentError("query response must contain exactly one non-empty query") |
| return query.strip()[:1000] |
|
|
|
|
| def identity_for(root: Path, experiment: Any, task: TaskSpec, harness: HarnessSpec, revision: str) -> RunIdentity: |
| model = load_models(root)[experiment.model_ids[0]] |
| return RunIdentity( |
| experiment_id=experiment.experiment_id, |
| task_id=task.task_id, |
| harness_id=harness.harness_id, |
| harness_hash=harness.config_hash, |
| model_id=model.model_id, |
| model_key=model.expected_inference_key, |
| model_config_hash=model.config_hash, |
| context_budget=experiment.context_budgets[0], |
| seed=experiment.seeds[0], |
| repetition=0, |
| repository_sha=task.base_commit, |
| code_revision=revision, |
| ) |
|
|
|
|
| def run_one_shot_and_queries( |
| root: Path, |
| repository: Path, |
| phase: str, |
| task_filter: set[str] | None = None, |
| harness_filter: set[str] | None = None, |
| ) -> dict[str, Any]: |
| revision = research_code_revision(root) |
| experiment = load_experiments(root)["E02"] |
| harness_catalog = load_harnesses(root) |
| model = load_models(root)[experiment.model_ids[0]] |
| tasks = selected_tasks(root, experiment, task_filter) |
| wanted = ( |
| {"H008", "H011", "H013", "H014", "H015"} |
| if phase == "one_shot" |
| else {"H010", "H012"} |
| ) |
| if harness_filter is not None: |
| wanted &= harness_filter |
| harnesses = [harness_catalog[item] for item in experiment.harness_ids if item in wanted] |
| if not tasks or not harnesses: |
| raise InteractiveExperimentError("filters selected no E02 model cells") |
| client = LMStudioClient(model, timeout_seconds=experiment.timeout_seconds) |
| discovery, resolved = client.resolve() |
| loaded = exclusive_loaded(discovery.native_models) |
| if loaded != (model.expected_inference_key,): |
| raise InteractiveExperimentError(f"E02 model phase requires exclusive Qwen residency: {loaded}") |
| tokenizer = QwenTokenCounter() |
| snapshot = GitSnapshot(repository) |
| rows: list[dict[str, Any]] = [] |
| for task in tasks: |
| for harness in harnesses: |
| context, allowed_paths, context_tokens = initial_context( |
| root, task, harness, snapshot, tokenizer |
| ) |
| prompt = localization_prompt(task, context) |
| if phase == "query": |
| stage = root / "results" / "staging" / "E02" / revision / harness.harness_id / task.task_id |
| stage.mkdir(parents=True, exist_ok=True) |
| output = stage / "query_stage.json" |
| if output.exists(): |
| rows.append(json.loads(output.read_text(encoding="utf-8"))) |
| continue |
| started = time.monotonic() |
| response = client.chat_completions( |
| resolved.inference_key, |
| [ |
| {"role": "system", "content": QUERY_SYSTEM}, |
| {"role": "user", "content": prompt}, |
| ], |
| max_tokens=1536, |
| ) |
| elapsed = time.monotonic() - started |
| query_violation: str | None = None |
| try: |
| query = parse_query(response) |
| query_source = "model" |
| except InteractiveExperimentError as exc: |
| query = task.statement |
| query_source = "issue_fallback" |
| query_violation = str(exc) |
| value = { |
| "task_id": task.task_id, |
| "harness_id": harness.harness_id, |
| "query": query, |
| "query_source": query_source, |
| "protocol_violation": query_violation, |
| "allowed_initial_paths": allowed_paths, |
| "context_tokens": context_tokens, |
| "prompt_sha256": sha256(prompt.encode()).hexdigest(), |
| "elapsed_seconds": elapsed, |
| "usage": response.get("usage", {}), |
| "response": response, |
| "prompt": prompt, |
| } |
| output.write_text(json.dumps(value, indent=2) + "\n", encoding="utf-8") |
| rows.append({key: value[key] for key in value if key not in {"response", "prompt"}}) |
| continue |
|
|
| identity = identity_for(root, experiment, task, harness, revision) |
| directory = run_directory(root / "results", identity) |
| if directory.exists(): |
| final = json.loads((directory / "final_metrics.json").read_text(encoding="utf-8")) |
| rows.append(final) |
| continue |
| started = time.monotonic() |
| response = client.chat_completions( |
| resolved.inference_key, |
| [ |
| {"role": "system", "content": LOCALIZE_SYSTEM}, |
| {"role": "user", "content": prompt}, |
| ], |
| max_tokens=model.max_tokens, |
| ) |
| elapsed = time.monotonic() - started |
| selection, protocol_violation = score_selection_response(response, set(allowed_paths)) |
| metrics = retrieval_metrics(selection["files"], task.gold_files) |
| final = { |
| "run_id": identity.run_id, |
| "experiment_id": "E02", |
| "task_id": task.task_id, |
| "harness_id": harness.harness_id, |
| "selected_files": selection["files"], |
| "reasoning": selection["reasoning"], |
| "metrics": metrics, |
| "context_tokens": context_tokens, |
| "prompt_sha256": sha256(prompt.encode()).hexdigest(), |
| "elapsed_seconds": elapsed, |
| "usage": response.get("usage", {}), |
| "model_calls": 1, |
| "protocol_violation": protocol_violation, |
| } |
| with EventWriter(root / "results", identity, asdict(harness), resolved.to_dict()) as writer: |
| writer.emit("run_started", {"confirmatory": True, "phase": phase}) |
| writer.emit("model_call", {"elapsed_seconds": elapsed, "usage": response.get("usage", {})}) |
| writer.write_artifact("prompt.txt", prompt) |
| writer.write_artifact("model_response.json", json.dumps(response, indent=2) + "\n") |
| writer.write_artifact( |
| "selection.json", |
| json.dumps({**selection, "protocol_violation": protocol_violation}, indent=2) + "\n", |
| ) |
| writer.write_artifact("final_metrics.json", json.dumps(final, indent=2) + "\n") |
| writer.emit( |
| "run_finished", |
| { |
| "status": "completed_with_protocol_violation" if protocol_violation else "completed", |
| "metrics": metrics, |
| "protocol_violation": protocol_violation, |
| }, |
| ) |
| rows.append(final) |
| return {"phase": phase, "code_revision": revision, "run_count": len(rows), "rows": rows} |
|
|
|
|
| def run_refined_retrieval( |
| root: Path, |
| repository: Path, |
| task_filter: set[str] | None = None, |
| harness_filter: set[str] | None = None, |
| ) -> dict[str, Any]: |
| revision = research_code_revision(root) |
| experiment = load_experiments(root)["E02"] |
| harness_catalog = load_harnesses(root) |
| embedding = load_embeddings(root)[experiment.embedding_id] |
| tasks = selected_tasks(root, experiment, task_filter) |
| wanted = {"H010", "H012"} if harness_filter is None else {"H010", "H012"} & harness_filter |
| harnesses = [harness_catalog[item] for item in experiment.harness_ids if item in wanted] |
| client = LMStudioEmbeddingClient(embedding, timeout_seconds=120.0) |
| runtime = client.resolve() |
| loaded = client.loaded_model_keys() |
| if tuple(loaded) != (embedding.model_key,): |
| raise InteractiveExperimentError(f"refined retrieval requires exclusive embedding residency: {loaded}") |
| snapshot = GitSnapshot(repository) |
| cache_path = root / "indexes" / "embeddings" / f"{embedding.config_hash}.sqlite3" |
| rows: list[dict[str, Any]] = [] |
| with SQLiteEmbeddingCache(cache_path, embedding) as cache: |
| for task in tasks: |
| chunks = chunk_snapshot(snapshot, task.base_commit, embedding.chunk_lines, embedding.chunk_overlap_lines, embedding.chunk_char_limit) |
| symbols = parse_snapshot(snapshot, task.base_commit) |
| exact = ExactRetriever(chunks) |
| lexical = BM25FuzzyRetriever(chunks) |
| syntax = SyntaxRetriever(symbols) |
| dense_base, _ = DenseRetriever.build(chunks, embedding, client, cache) |
| dense = FaissFlatRetriever(dense_base) |
| graph = SymbolGraph(symbols) |
| for harness in harnesses: |
| stage = root / "results" / "staging" / "E02" / revision / harness.harness_id / task.task_id |
| query_stage = stage / "query_stage.json" |
| if not query_stage.exists(): |
| raise InteractiveExperimentError(f"missing query stage: {query_stage}") |
| output = stage / "refined_ranking.json" |
| if output.exists(): |
| rows.append({"task_id": task.task_id, "harness_id": harness.harness_id, "status": "resumed"}) |
| continue |
| query = json.loads(query_stage.read_text(encoding="utf-8"))["query"] |
| ranking = retrieve_treatment(harness, query, exact, lexical, syntax, dense, graph, 200) |
| payload = [ |
| { |
| "rank": rank, |
| "path": item.path, |
| "line_start": item.line_start, |
| "line_end": item.line_end, |
| "score": item.score, |
| "source": item.source, |
| "symbol": item.symbol, |
| } |
| for rank, item in enumerate(ranking, start=1) |
| ] |
| output.write_text(json.dumps(payload, indent=2) + "\n", encoding="utf-8") |
| rows.append({"task_id": task.task_id, "harness_id": harness.harness_id, "query": query, "candidate_count": len(payload)}) |
| return {"phase": "refined_retrieval", "embedding_runtime": runtime, "run_count": len(rows), "rows": rows} |
|
|
|
|
| def run_iterative_final( |
| root: Path, |
| repository: Path, |
| task_filter: set[str] | None = None, |
| harness_filter: set[str] | None = None, |
| ) -> dict[str, Any]: |
| revision = research_code_revision(root) |
| experiment = load_experiments(root)["E02"] |
| model = load_models(root)[experiment.model_ids[0]] |
| harness_catalog = load_harnesses(root) |
| tasks = selected_tasks(root, experiment, task_filter) |
| wanted = {"H010", "H012"} if harness_filter is None else {"H010", "H012"} & harness_filter |
| harnesses = [harness_catalog[item] for item in experiment.harness_ids if item in wanted] |
| client = LMStudioClient(model, timeout_seconds=experiment.timeout_seconds) |
| discovery, resolved = client.resolve() |
| loaded = exclusive_loaded(discovery.native_models) |
| if loaded != (model.expected_inference_key,): |
| raise InteractiveExperimentError(f"iterative final phase requires exclusive Qwen residency: {loaded}") |
| tokenizer = QwenTokenCounter() |
| snapshot = GitSnapshot(repository) |
| rows: list[dict[str, Any]] = [] |
| for task in tasks: |
| for harness in harnesses: |
| identity = identity_for(root, experiment, task, harness, revision) |
| directory = run_directory(root / "results", identity) |
| if directory.exists(): |
| rows.append(json.loads((directory / "final_metrics.json").read_text(encoding="utf-8"))) |
| continue |
| stage = root / "results" / "staging" / "E02" / revision / harness.harness_id / task.task_id |
| query_value = json.loads((stage / "query_stage.json").read_text(encoding="utf-8")) |
| refined_records = load_ranking_records(stage / "refined_ranking.json") |
| refined_candidates = candidates_from_records(snapshot, task.base_commit, refined_records) |
| refined_text, refined_paths, refined_tokens = pack_snippets( |
| tokenizer, refined_candidates, EVIDENCE_TOKEN_BUDGET // 2 |
| ) |
| initial_text, initial_paths, initial_tokens = initial_context(root, task, harness, snapshot, tokenizer) |
| |
| if harness.interface == "specialized": |
| initial_text, initial_paths, initial_tokens = pack_specialized_channels( |
| root, task.task_id, snapshot, task.base_commit, tokenizer, EVIDENCE_TOKEN_BUDGET // 2 |
| ) |
| else: |
| initial_records = load_ranking_records(ranking_path(root, "E01", "H008", task.task_id)) |
| initial_candidates = candidates_from_records(snapshot, task.base_commit, initial_records) |
| initial_text, initial_paths, initial_tokens = pack_snippets( |
| tokenizer, initial_candidates, EVIDENCE_TOKEN_BUDGET // 2 |
| ) |
| context = ( |
| f"FOLLOW-UP QUERY:\n{query_value['query']}\n\nINITIAL EVIDENCE:\n{initial_text}" |
| f"\n\nREFINED EVIDENCE:\n{refined_text}" |
| ) |
| prompt = localization_prompt(task, context, phase="iterative") |
| allowed = tuple(dict.fromkeys((*initial_paths, *refined_paths))) |
| started = time.monotonic() |
| response = client.chat_completions( |
| resolved.inference_key, |
| [{"role": "system", "content": LOCALIZE_SYSTEM}, {"role": "user", "content": prompt}], |
| max_tokens=model.max_tokens, |
| ) |
| elapsed = time.monotonic() - started |
| selection, protocol_violation = score_selection_response(response, set(allowed)) |
| metrics = retrieval_metrics(selection["files"], task.gold_files) |
| first_usage = query_value.get("usage", {}) |
| second_usage = response.get("usage", {}) |
| final = { |
| "run_id": identity.run_id, |
| "experiment_id": "E02", |
| "task_id": task.task_id, |
| "harness_id": harness.harness_id, |
| "selected_files": selection["files"], |
| "reasoning": selection["reasoning"], |
| "query": query_value["query"], |
| "metrics": metrics, |
| "context_tokens": initial_tokens + refined_tokens, |
| "prompt_sha256": sha256(prompt.encode()).hexdigest(), |
| "elapsed_seconds": query_value["elapsed_seconds"] + elapsed, |
| "usage": {"first": first_usage, "second": second_usage}, |
| "model_calls": 2, |
| "protocol_violation": protocol_violation, |
| } |
| with EventWriter(root / "results", identity, asdict(harness), resolved.to_dict()) as writer: |
| writer.emit("run_started", {"confirmatory": True, "phase": "iterative"}) |
| writer.emit("model_call", {"stage": "query", "elapsed_seconds": query_value["elapsed_seconds"], "usage": first_usage}) |
| writer.emit("model_call", {"stage": "final", "elapsed_seconds": elapsed, "usage": second_usage}) |
| writer.write_artifact("prompt_query.txt", query_value["prompt"]) |
| writer.write_artifact("prompt_final.txt", prompt) |
| writer.write_artifact("response_query.json", json.dumps(query_value["response"], indent=2) + "\n") |
| writer.write_artifact("response_final.json", json.dumps(response, indent=2) + "\n") |
| writer.write_artifact( |
| "selection.json", |
| json.dumps({**selection, "protocol_violation": protocol_violation}, indent=2) + "\n", |
| ) |
| writer.write_artifact("final_metrics.json", json.dumps(final, indent=2) + "\n") |
| writer.emit( |
| "run_finished", |
| { |
| "status": "completed_with_protocol_violation" if protocol_violation else "completed", |
| "metrics": metrics, |
| "protocol_violation": protocol_violation, |
| }, |
| ) |
| rows.append(final) |
| return {"phase": "iterative_final", "code_revision": revision, "run_count": len(rows), "rows": rows} |
|
|