evalstate/diffusers-pr-api
0
1from __future__ import annotations2 3import asyncio4import copy5import json6import os7import re8import shutil9import sys10from collections import Counter, defaultdict11from dataclasses import dataclass12from datetime import UTC, datetime13from pathlib import Path14from typing import Any15 16from pydantic import BaseModel, Field17from rank_bm25 import BM25Okapi18 19from slop_farmer.config import AnalysisOptions, MarkdownReportOptions20from slop_farmer.data.links import build_text_link_rows21from slop_farmer.data.parquet_io import read_json, read_parquet_rows, write_text22from slop_farmer.data.snapshot_source import resolve_snapshot_source_dir23from slop_farmer.reports.analysis_cache import (24 HYBRID_REVIEW_CACHE_SCHEMA_VERSION,25 PREPARED_REVIEW_UNIT_SCHEMA_VERSION,26 HybridReviewCacheEntry,27 HybridReviewCacheKey,28 HybridReviewCacheManifest,29 HybridReviewCacheStore,30 HybridReviewSettingsFingerprint,31 build_hybrid_review_cache_key,32 hybrid_review_cache_dir,33)34from slop_farmer.reports.pr_heuristics import (35 build_template_cleanup_settings,36 compile_cluster_suppression_rules,37 strip_pull_request_template,38 suppressed_pull_request_reasons,39)40 41LINK_KEY_FIELDS = (42 "repo",43 "source_type",44 "source_number",45 "source_github_id",46 "target_owner",47 "target_repo",48 "target_number",49 "link_type",50 "link_origin",51)52STOPWORDS = {53 "a",54 "an",55 "and",56 "are",57 "as",58 "at",59 "be",60 "by",61 "for",62 "from",63 "how",64 "if",65 "in",66 "into",67 "is",68 "it",69 "of",70 "on",71 "or",72 "that",73 "the",74 "this",75 "to",76 "was",77 "were",78 "with",79}80TOKEN_PATTERN = re.compile(r"[a-z0-9_]+")81HUNK_HEADER_PATTERN = re.compile(r"^@@ -\d+(?:,\d+)? \+(?P<start>\d+)(?:,(?P<count>\d+))? @@")82LLM_PROVIDER_ENV_VARS = (83 "OPENAI_API_KEY",84 "ANTHROPIC_API_KEY",85 "GOOGLE_API_KEY",86 "DEEPSEEK_API_KEY",87)88LLM_PACKET_CHARS_PER_TOKEN = 489LLM_MAX_INPUT_TOKENS = 60_00090LLM_MAX_NODES_PER_PACKET = 4891LLM_MAX_SOFT_PAIRS_PER_PACKET = 7292LLM_MAX_DIFF_CHARS_PER_ITEM = 1_20093LLM_MAX_FILENAMES_PER_ITEM = 1694LLM_SKIP_EVALUATOR_ABOVE_TOKENS = 60_00095LLM_OVERFLOW_POLICY = "truncate_then_skip"96LLM_SHARED_TARGET_MAX_NEIGHBORS_PER_PR = 397LLM_SHARED_TARGET_MAX_EXTRA_PAIRS_PER_TARGET = 1898LLM_SHARED_TARGET_MIN_TEXT_JACCARD = 0.199CLUSTER_ANALYST_PROMPT_VERSION = "1.0"100CLUSTER_EVALUATOR_PROMPT_VERSION = "1.0"101CLUSTER_ANALYST_INSTRUCTION = (102 "You analyze clustered GitHub issues and pull requests for duplicate triage. "103 "Return a short summary, confidence between 0 and 1, concise reasons for canonical issue/PR choices, "104 "concise reasons for global best issue/PR suitability, and accept/reject verdicts for each soft edge candidate. "105 "Only accept a soft edge when the two artifacts look like the same underlying bug or change. "106 "Use titles, descriptions, explicit issue targets, changed filenames, and diff previews when available. "107 "For pull requests, be strict: accept only when the PRs appear to fix the same concrete code-path problem and could plausibly be merged into one PR. "108 "Do not merge PRs just because they mention the same tracking issue, touch the same broad subsystem, or both change documentation/tests."109)110CLUSTER_EVALUATOR_INSTRUCTION = (111 "You review the analyst output for precision. Accept only when the summary is grounded in the packet "112 "and every soft-edge verdict is conservative. Reject if the analyst overstates evidence. "113 "For pull-request pairs, reject if the two changes do not look mergeable into a single PR for the same bugfix."114)115 116 117class SoftEdgeVerdict(BaseModel):118 left: str119 right: str120 accept: bool121 reason: str122 123 124class ClusterAnalystResponse(BaseModel):125 summary: str126 confidence: float127 canonical_issue_reason: str | None = None128 canonical_pr_reason: str | None = None129 best_issue_reason: str | None = None130 best_pr_reason: str | None = None131 soft_edge_verdicts: list[SoftEdgeVerdict] = Field(default_factory=list)132 133 134class ClusterEvaluatorResponse(BaseModel):135 accept: bool136 feedback: str = ""137 138 139class PrFileAreaEntry(BaseModel):140 filename: str141 left_ranges: list[list[int]]142 right_ranges: list[list[int]]143 144 145class PrComparisonEntry(BaseModel):146 left_pr_number: int147 right_pr_number: int148 code_similarity: float149 size_similarity: float150 file_overlap: float151 area_overlap: float152 patch_similarity: float153 shared_filenames: list[str]154 shared_file_areas: list[PrFileAreaEntry]155 156 157class MetaBugEntry(BaseModel):158 cluster_id: str159 summary: str160 status: str161 confidence: float162 canonical_issue_number: int | None163 canonical_pr_number: int | None164 issue_numbers: list[int]165 pr_numbers: list[int]166 evidence_types: list[str]167 pr_comparisons: list[PrComparisonEntry] = Field(default_factory=list)168 169 170class DuplicateIssuesEntry(BaseModel):171 cluster_id: str172 canonical_issue_number: int173 duplicate_issue_numbers: list[int]174 reason: str175 176 177class DuplicatePrsEntry(BaseModel):178 cluster_id: str179 canonical_pr_number: int180 duplicate_pr_numbers: list[int]181 target_issue_number: int | None182 reason: str183 184 185class BestIssueEntry(BaseModel):186 cluster_id: str187 issue_number: int188 reason: str189 score: float190 191 192class BestPrEntry(BaseModel):193 cluster_id: str194 pr_number: int195 reason: str196 score: float197 198 199class AnalysisReport(BaseModel):200 schema_version: str201 repo: str202 snapshot_id: str203 generated_at: str204 evidence_quality: str205 llm_enrichment: bool206 meta_bugs: list[MetaBugEntry]207 duplicate_issues: list[DuplicateIssuesEntry]208 duplicate_prs: list[DuplicatePrsEntry]209 best_issue: BestIssueEntry | None210 best_pr: BestPrEntry | None211 212 213@dataclass(slots=True)214class SnapshotData:215 repo: str216 snapshot_id: str217 snapshot_dir: Path218 manifest: dict[str, Any]219 issues: list[dict[str, Any]]220 pull_requests: list[dict[str, Any]]221 comments: list[dict[str, Any]]222 reviews: list[dict[str, Any]]223 review_comments: list[dict[str, Any]]224 pr_files: list[dict[str, Any]]225 pr_diffs: list[dict[str, Any]]226 links: list[dict[str, Any]]227 events: list[dict[str, Any]]228 evidence_quality: str229 230 231@dataclass(slots=True)232class ArtifactFeature:233 node_id: str234 kind: str235 number: int236 row: dict[str, Any]237 tokens: list[str]238 title_tokens: set[str]239 title_length: int240 body_length: int241 discussion_activity: int242 review_activity: int243 inbound_references: int244 explicit_issue_links: int245 explicit_issue_targets: list[int]246 diff_size: int247 filenames: list[str]248 diff_preview: str | None249 file_ranges_by_name: dict[str, list[tuple[int, int]]]250 patch_tokens: list[str]251 252 253@dataclass(slots=True)254class ClusterRecord:255 cluster_id: str256 nodes: list[str]257 issue_numbers: list[int]258 pr_numbers: list[int]259 evidence_types: list[str]260 canonical_issue_number: int | None261 canonical_pr_number: int | None262 target_issue_number: int | None263 summary: str264 status: str265 confidence: float266 canonical_issue_reason: str | None267 canonical_pr_reason: str | None268 best_issue_reason: str | None269 best_pr_reason: str | None270 cluster_score: float271 best_issue_score: float | None272 best_pr_score: float | None273 274 275@dataclass(frozen=True, slots=True)276class PacketBudget:277 node_count: int278 item_count: int279 soft_pair_count: int280 serialized_chars: int281 estimated_input_tokens: int282 estimated_eval_tokens: int283 284 285@dataclass(frozen=True, slots=True)286class PreparedLlmPacket:287 packet: dict[str, Any]288 budget: PacketBudget289 original_budget: PacketBudget290 trimmed: bool291 aggressively_trimmed: bool292 split: bool293 294 295@dataclass(frozen=True, slots=True)296class ClusterAnalysisCallResult:297 analyst_result: ClusterAnalystResponse | None298 evaluator_result: ClusterEvaluatorResponse | None299 error_kind: str | None300 error_message: str | None301 evaluator_used: bool302 retried: bool303 304 305@dataclass(frozen=True, slots=True)306class AnalysisBuildResult:307 report: AnalysisReport308 llm_reviews: list[dict[str, Any]]309 310 311@dataclass(frozen=True, slots=True)312class SoftPairReviewUnitMeta:313 label: str314 component_index: int315 component_count: int316 review_unit_index: int317 review_unit_count: int318 cluster_id: str319 prefix: str320 nodes: tuple[str, ...]321 soft_pairs: tuple[str, ...]322 component_budget: PacketBudget323 budget: PacketBudget324 prepared_review_unit_hash: str | None325 trimmed: bool326 aggressively_trimmed: bool327 split: bool328 329 330@dataclass(frozen=True, slots=True)331class PendingSoftPairReview:332 meta: SoftPairReviewUnitMeta333 prepared: PreparedLlmPacket334 cache_key: HybridReviewCacheKey335 336 337@dataclass(frozen=True, slots=True)338class CompletedSoftPairReview:339 meta: SoftPairReviewUnitMeta340 result: ClusterAnalysisCallResult | None341 status: str342 reason: str | None343 source: str | None344 cache_hit: bool345 346 347def _hybrid_review_cache_manifest() -> HybridReviewCacheManifest:348 return HybridReviewCacheManifest(349 cache_schema_version=HYBRID_REVIEW_CACHE_SCHEMA_VERSION,350 prepared_review_unit_schema_version=PREPARED_REVIEW_UNIT_SCHEMA_VERSION,351 analyst_prompt_version=CLUSTER_ANALYST_PROMPT_VERSION,352 evaluator_prompt_version=CLUSTER_EVALUATOR_PROMPT_VERSION,353 hybrid_review_settings=HybridReviewSettingsFingerprint(354 llm_max_input_tokens=LLM_MAX_INPUT_TOKENS,355 llm_max_nodes_per_packet=LLM_MAX_NODES_PER_PACKET,356 llm_max_soft_pairs_per_packet=LLM_MAX_SOFT_PAIRS_PER_PACKET,357 llm_max_diff_chars_per_item=LLM_MAX_DIFF_CHARS_PER_ITEM,358 llm_max_filenames_per_item=LLM_MAX_FILENAMES_PER_ITEM,359 llm_skip_evaluator_above_tokens=LLM_SKIP_EVALUATOR_ABOVE_TOKENS,360 llm_overflow_policy=LLM_OVERFLOW_POLICY,361 ),362 )363 364 365def _prepared_review_unit_payload(prepared: PreparedLlmPacket) -> dict[str, Any]:366 return {367 "packet": copy.deepcopy(prepared.packet),368 "budget": _packet_budget_json(prepared.budget),369 "original_budget": _packet_budget_json(prepared.original_budget),370 "trimmed": prepared.trimmed,371 "aggressively_trimmed": prepared.aggressively_trimmed,372 "split": prepared.split,373 }374 375 376def _cluster_analysis_call_result_payload(result: ClusterAnalysisCallResult) -> dict[str, Any]:377 return {378 "analyst_result": (379 None if result.analyst_result is None else result.analyst_result.model_dump(mode="json")380 ),381 "evaluator_result": (382 None383 if result.evaluator_result is None384 else result.evaluator_result.model_dump(mode="json")385 ),386 "error_kind": result.error_kind,387 "error_message": result.error_message,388 "evaluator_used": result.evaluator_used,389 "retried": result.retried,390 }391 392 393def _cluster_analysis_call_result_from_payload(394 payload: dict[str, Any],395) -> ClusterAnalysisCallResult:396 return ClusterAnalysisCallResult(397 analyst_result=(398 None399 if payload.get("analyst_result") is None400 else ClusterAnalystResponse.model_validate(payload["analyst_result"])401 ),402 evaluator_result=(403 None404 if payload.get("evaluator_result") is None405 else ClusterEvaluatorResponse.model_validate(payload["evaluator_result"])406 ),407 error_kind=payload.get("error_kind"),408 error_message=payload.get("error_message"),409 evaluator_used=bool(payload.get("evaluator_used", False)),410 retried=bool(payload.get("retried", False)),411 )412 413 414def _cacheable_cluster_analysis_result(result: ClusterAnalysisCallResult) -> bool:415 return result.analyst_result is not None and result.error_kind is None416 417 418def run_analysis(options: AnalysisOptions) -> Path:419 if options.snapshot_dir is not None and options.hf_repo_id:420 raise ValueError("--snapshot-dir and --hf-repo-id are mutually exclusive")421 warning = _llm_fallback_warning(options)422 if warning:423 _analysis_log(warning)424 snapshot_dir = _resolve_snapshot_dir(options)425 snapshot = _load_snapshot(snapshot_dir)426 _maybe_carry_forward_hybrid_review_cache(snapshot, enabled=options.cached_analysis)427 build = asyncio.run(_build_report(snapshot, options))428 output_path = options.output or (snapshot_dir / "analysis-report.json")429 output_path.parent.mkdir(parents=True, exist_ok=True)430 write_text(json.dumps(build.report.model_dump(mode="json"), indent=2) + "\n", output_path)431 llm_reviews_path = _llm_reviews_output_path(output_path)432 if build.llm_reviews:433 write_text(434 json.dumps(435 {436 "schema_version": "1.0",437 "repo": build.report.repo,438 "snapshot_id": build.report.snapshot_id,439 "generated_at": build.report.generated_at,440 "model": options.model,441 "reviews": build.llm_reviews,442 },443 indent=2,444 )445 + "\n",446 llm_reviews_path,447 )448 elif llm_reviews_path.exists():449 llm_reviews_path.unlink()450 _log_hybrid_review_cache_summary(build.llm_reviews, enabled=options.cached_analysis)451 return output_path452 453 454def _analysis_log(message: str) -> None:455 stamp = datetime.now(tz=UTC).strftime("%H:%M:%SZ")456 print(f"[{stamp}] {message}", file=sys.stderr, flush=True)457 458 459def _llm_reviews_output_path(output_path: Path) -> Path:460 return output_path.with_name(f"{output_path.stem}.llm-reviews.json")461 462 463def _llm_fallback_warning(options: AnalysisOptions) -> str | None:464 if options.ranking_backend != "hybrid":465 return None466 if _can_use_fast_agent():467 return None468 return (469 "Analyze requested ranking-backend=hybrid but fast-agent LLM enrichment is unavailable; "470 "reusing cached hybrid review results when available and falling back to deterministic-only clustering "471 "for cache misses. "472 "Install the llm extra and set one of "473 f"{', '.join(LLM_PROVIDER_ENV_VARS)}."474 )475 476 477def _maybe_carry_forward_hybrid_review_cache(snapshot: SnapshotData, *, enabled: bool) -> None:478 if not enabled:479 return480 current_cache_dir = hybrid_review_cache_dir(snapshot.snapshot_dir)481 if current_cache_dir.exists():482 _analysis_log(483 f"Cached analysis enabled: using existing analysis-state in {current_cache_dir}"484 )485 return486 watermark = snapshot.manifest.get("watermark")487 if not isinstance(watermark, dict):488 _analysis_log("Cached analysis enabled: no previous snapshot recorded; starting fresh")489 return490 previous_snapshot_dir = watermark.get("previous_snapshot_dir")491 if not isinstance(previous_snapshot_dir, str) or not previous_snapshot_dir:492 _analysis_log("Cached analysis enabled: no previous snapshot recorded; starting fresh")493 return494 previous_cache_dir = hybrid_review_cache_dir(Path(previous_snapshot_dir))495 if not previous_cache_dir.exists():496 _analysis_log(497 "Cached analysis enabled: previous snapshot has no analysis-state; starting fresh"498 )499 return500 shutil.copytree(previous_cache_dir, current_cache_dir)501 _analysis_log(502 f"Cached analysis enabled: copied analysis-state from {previous_cache_dir} to {current_cache_dir}"503 )504 505 506def _log_hybrid_review_cache_summary(llm_reviews: list[dict[str, Any]], *, enabled: bool) -> None:507 if not enabled:508 return509 if not llm_reviews:510 _analysis_log("Hybrid review cache summary: no LLM review units were produced")511 return512 reviewed = [review for review in llm_reviews if review.get("status") == "reviewed"]513 cache_hits = [review for review in reviewed if review.get("cache_hit")]514 cache_sourced = [review for review in reviewed if review.get("source") == "cache"]515 llm_sourced = [review for review in reviewed if review.get("source") == "llm"]516 skipped = [review for review in llm_reviews if review.get("status") != "reviewed"]517 hit_rate = 100.0 * len(cache_hits) / len(reviewed) if reviewed else 0.0518 _analysis_log(519 "Hybrid review cache summary: "520 f"{len(cache_hits)}/{len(reviewed)} reviewed units reused from cache "521 f"({hit_rate:.1f}%); "522 f"source_cache={len(cache_sourced)}, source_llm={len(llm_sourced)}, skipped={len(skipped)}"523 )524 if skipped:525 reasons = Counter(str(review.get("reason")) for review in skipped if review.get("reason"))526 if reasons:527 formatted = ", ".join(f"{reason}={count}" for reason, count in reasons.most_common(5))528 _analysis_log(f"Hybrid review cache skipped reasons: {formatted}")529 530 531def render_markdown_report(options: MarkdownReportOptions) -> Path:532 input_path = options.input.resolve()533 report = AnalysisReport.model_validate(read_json(input_path))534 snapshot_dir = _resolve_markdown_snapshot_dir(input_path, options.snapshot_dir)535 issue_map, pr_map = _report_artifact_maps(snapshot_dir)536 output_path = (options.output or input_path.with_suffix(".md")).resolve()537 markdown = _markdown_report_text(538 report=report,539 issue_map=issue_map,540 pr_map=pr_map,541 )542 write_text(markdown, output_path)543 return output_path544 545 546def _resolve_markdown_snapshot_dir(input_path: Path, snapshot_dir: Path | None) -> Path | None:547 if snapshot_dir is not None:548 return snapshot_dir.resolve()549 candidate = input_path.parent.resolve()550 if (candidate / "issues.parquet").exists() or (candidate / "pull_requests.parquet").exists():551 return candidate552 return None553 554 555def _report_artifact_maps(556 snapshot_dir: Path | None,557) -> tuple[dict[int, dict[str, Any]], dict[int, dict[str, Any]]]:558 if snapshot_dir is None:559 return {}, {}560 issues = {561 int(row["number"]): row562 for row in read_parquet_rows(snapshot_dir / "issues.parquet")563 if row.get("number") is not None564 }565 pull_requests = {566 int(row["number"]): row567 for row in read_parquet_rows(snapshot_dir / "pull_requests.parquet")568 if row.get("number") is not None569 }570 return issues, pull_requests571 572 573def _markdown_report_text(574 *,575 report: AnalysisReport,576 issue_map: dict[int, dict[str, Any]],577 pr_map: dict[int, dict[str, Any]],578) -> str:579 lines = [580 f"# Analysis Report: {report.repo}",581 "",582 f"- Snapshot: `{report.snapshot_id}`",583 f"- Generated: `{report.generated_at}`",584 f"- Evidence quality: `{report.evidence_quality}`",585 f"- LLM enrichment: `{str(report.llm_enrichment).lower()}`",586 f"- Meta bugs: `{len(report.meta_bugs)}`",587 ]588 if report.best_issue is not None:589 lines.append(590 f"- Best issue: {_artifact_markdown_link(report.repo, 'issue', report.best_issue.issue_number, issue_map.get(report.best_issue.issue_number))}"591 )592 if report.best_pr is not None:593 lines.append(594 f"- Best PR: {_artifact_markdown_link(report.repo, 'pull_request', report.best_pr.pr_number, pr_map.get(report.best_pr.pr_number))}"595 )596 lines.append("")597 598 ordered_meta_bugs = sorted(599 report.meta_bugs,600 key=lambda entry: _meta_bug_sort_key(entry, issue_map, pr_map),601 )602 if not ordered_meta_bugs:603 lines.append("No meta bugs found.")604 lines.append("")605 return "\n".join(lines)606 607 for meta_bug in ordered_meta_bugs:608 lines.extend(_meta_bug_markdown_lines(report.repo, meta_bug, issue_map, pr_map))609 return "\n".join(lines).rstrip() + "\n"610 611 612def _meta_bug_markdown_lines(613 repo: str,614 meta_bug: MetaBugEntry,615 issue_map: dict[int, dict[str, Any]],616 pr_map: dict[int, dict[str, Any]],617) -> list[str]:618 artifact_count = len(meta_bug.issue_numbers) + len(meta_bug.pr_numbers)619 latest_activity = _meta_bug_latest_activity(meta_bug, issue_map, pr_map)620 issue_numbers_to_render = [621 number for number in meta_bug.issue_numbers if number != meta_bug.canonical_issue_number622 ]623 lines = [624 f"## {meta_bug.summary}",625 "",626 f"- Cluster: `{meta_bug.cluster_id}`",627 f"- Status: `{meta_bug.status}`",628 f"- Confidence: `{meta_bug.confidence:.3f}`",629 f"- Artifacts: `{artifact_count}`",630 f"- Latest activity: `{latest_activity}`",631 ]632 if meta_bug.canonical_issue_number is not None:633 lines.append(634 f"- Canonical issue: {_artifact_markdown_link(repo, 'issue', meta_bug.canonical_issue_number, issue_map.get(meta_bug.canonical_issue_number))}"635 )636 if meta_bug.canonical_pr_number is not None:637 lines.append(638 f"- Canonical PR: {_artifact_markdown_link(repo, 'pull_request', meta_bug.canonical_pr_number, pr_map.get(meta_bug.canonical_pr_number))}"639 )640 if meta_bug.evidence_types:641 lines.append(f"- Evidence: `{', '.join(meta_bug.evidence_types)}`")642 lines.append("")643 644 if issue_numbers_to_render:645 lines.append("### Issues")646 lines.append("")647 for number in _sorted_artifact_numbers(issue_numbers_to_render, issue_map):648 lines.append(649 f"- {_artifact_markdown_link(repo, 'issue', number, issue_map.get(number))}{_artifact_suffix(issue_map.get(number), 'issue')}"650 )651 lines.append("")652 653 if meta_bug.pr_numbers:654 lines.append("### PRs")655 lines.append("")656 for number in _sorted_artifact_numbers(meta_bug.pr_numbers, pr_map):657 lines.append(658 f"- {_artifact_markdown_link(repo, 'pull_request', number, pr_map.get(number))}{_artifact_suffix(pr_map.get(number), 'pull_request')}"659 )660 lines.append("")661 662 if meta_bug.pr_comparisons:663 lines.append("### PR comparison")664 lines.append("")665 for comparison in meta_bug.pr_comparisons:666 shared_files = ", ".join(f"`{name}`" for name in comparison.shared_filenames) or "none"667 lines.append(668 f"- PR #{comparison.left_pr_number} vs PR #{comparison.right_pr_number}: "669 f"code `{comparison.code_similarity:.3f}`, "670 f"size `{comparison.size_similarity:.3f}`, "671 f"files `{comparison.file_overlap:.3f}`, "672 f"areas `{comparison.area_overlap:.3f}`, "673 f"patch `{comparison.patch_similarity:.3f}`; "674 f"shared files: {shared_files}"675 )676 lines.append("")677 678 return lines679 680 681def _meta_bug_sort_key(682 meta_bug: MetaBugEntry,683 issue_map: dict[int, dict[str, Any]],684 pr_map: dict[int, dict[str, Any]],685) -> tuple[int, float, int, str]:686 artifact_count = len(meta_bug.issue_numbers) + len(meta_bug.pr_numbers)687 latest_activity = _meta_bug_latest_activity_dt(meta_bug, issue_map, pr_map).timestamp()688 largest_number = max([*meta_bug.issue_numbers, *meta_bug.pr_numbers], default=0)689 return (-artifact_count, -latest_activity, -largest_number, meta_bug.cluster_id)690 691 692def _meta_bug_latest_activity(693 meta_bug: MetaBugEntry, issue_map: dict[int, dict[str, Any]], pr_map: dict[int, dict[str, Any]]694) -> str:695 latest_row = _meta_bug_latest_row(meta_bug, issue_map, pr_map)696 if latest_row is None:697 return "unknown"698 return str(699 latest_row.get("updated_at")700 or latest_row.get("created_at")701 or latest_row.get("closed_at")702 or "unknown"703 )704 705 706def _meta_bug_latest_activity_dt(707 meta_bug: MetaBugEntry,708 issue_map: dict[int, dict[str, Any]],709 pr_map: dict[int, dict[str, Any]],710) -> datetime:711 latest_row = _meta_bug_latest_row(meta_bug, issue_map, pr_map)712 if latest_row is None:713 return datetime(1970, 1, 1, tzinfo=UTC)714 return _row_activity_dt(latest_row)715 716 717def _meta_bug_latest_row(718 meta_bug: MetaBugEntry,719 issue_map: dict[int, dict[str, Any]],720 pr_map: dict[int, dict[str, Any]],721) -> dict[str, Any] | None:722 rows = [issue_map[number] for number in meta_bug.issue_numbers if number in issue_map]723 rows.extend(pr_map[number] for number in meta_bug.pr_numbers if number in pr_map)724 if not rows:725 return None726 return max(rows, key=_row_activity_dt)727 728 729def _sorted_artifact_numbers(numbers: list[int], row_map: dict[int, dict[str, Any]]) -> list[int]:730 return sorted(731 numbers,732 key=lambda number: (733 -_row_activity_dt(row_map.get(number)).timestamp(),734 -number,735 ),736 )737 738 739def _row_activity_dt(row: dict[str, Any] | None) -> datetime:740 if not row:741 return datetime(1970, 1, 1, tzinfo=UTC)742 for field in ("updated_at", "created_at", "closed_at", "merged_at"):743 value = row.get(field)744 if not value:745 continue746 try:747 return _parse_dt(str(value))748 except ValueError:749 continue750 return datetime(1970, 1, 1, tzinfo=UTC)751 752 753def _artifact_markdown_link(repo: str, kind: str, number: int, row: dict[str, Any] | None) -> str:754 title = _artifact_title(kind, number, row)755 url = _artifact_url(repo, kind, number, row)756 return f"[{title}]({url})"757 758 759def _artifact_title(kind: str, number: int, row: dict[str, Any] | None) -> str:760 prefix = "PR" if kind == "pull_request" else "Issue"761 title = str((row or {}).get("title") or "").strip()762 if not title and kind == "pull_request":763 body = str((row or {}).get("body") or "").strip()764 if body:765 title = body.splitlines()[0].strip()[:120]766 if title:767 return f"{prefix} #{number}: {title}"768 return f"{prefix} #{number}"769 770 771def _artifact_url(repo: str, kind: str, number: int, row: dict[str, Any] | None) -> str:772 html_url = str((row or {}).get("html_url") or "").strip()773 if html_url:774 return html_url775 if repo:776 path = "pull" if kind == "pull_request" else "issues"777 return f"https://github.com/{repo}/{path}/{number}"778 return "#"779 780 781def _artifact_suffix(row: dict[str, Any] | None, kind: str) -> str:782 if not row:783 return ""784 details: list[str] = []785 state = str(row.get("state") or "").strip()786 if state:787 details.append(state)788 if kind == "pull_request":789 if bool(row.get("merged")):790 details.append("merged")791 if bool(row.get("draft")):792 details.append("draft")793 timestamp = row.get("updated_at") or row.get("created_at")794 if timestamp:795 details.append(str(timestamp))796 if not details:797 return ""798 return f" ({', '.join(details)})"799 800 801def _resolve_snapshot_dir(options: AnalysisOptions) -> Path:802 return resolve_snapshot_source_dir(803 snapshot_dir=options.snapshot_dir,804 local_snapshots_root=options.output_dir.resolve() / "snapshots",805 hf_repo_id=options.hf_repo_id,806 hf_revision=options.hf_revision,807 hf_materialize_dir=options.hf_materialize_dir,808 hf_output_dir=options.output_dir,809 )810 811 812def _load_snapshot(snapshot_dir: Path) -> SnapshotData:813 manifest_path = snapshot_dir / "manifest.json"814 manifest = read_json(manifest_path) if manifest_path.exists() else {}815 816 issues = read_parquet_rows(snapshot_dir / "issues.parquet")817 pull_requests = read_parquet_rows(snapshot_dir / "pull_requests.parquet")818 comments = read_parquet_rows(snapshot_dir / "comments.parquet")819 reviews = read_parquet_rows(snapshot_dir / "reviews.parquet")820 review_comments = read_parquet_rows(snapshot_dir / "review_comments.parquet")821 pr_files = read_parquet_rows(snapshot_dir / "pr_files.parquet")822 pr_diffs = read_parquet_rows(snapshot_dir / "pr_diffs.parquet")823 links = read_parquet_rows(snapshot_dir / "links.parquet")824 events_path = snapshot_dir / "events.parquet"825 events = read_parquet_rows(events_path)826 if not any(827 [828 issues,829 pull_requests,830 comments,831 reviews,832 review_comments,833 pr_files,834 pr_diffs,835 links,836 events,837 ]838 ):839 parquet_files = sorted(str(path.name) for path in snapshot_dir.glob("*.parquet"))840 raise FileNotFoundError(841 f"No analysis tables found in {snapshot_dir}. "842 f"Expected local files like issues.parquet/pull_requests.parquet. "843 f"Found parquet files: {parquet_files or 'none'}. "844 "Use --hf-repo-id for Hugging Face datasets or point --snapshot-dir at a local slop-farmer snapshot."845 )846 847 repo = (848 manifest.get("repo")849 or (issues[0]["repo"] if issues else None)850 or (pull_requests[0]["repo"] if pull_requests else None)851 or (comments[0]["repo"] if comments else None)852 or ""853 )854 snapshot_id = manifest.get("snapshot_id") or snapshot_dir.name855 evidence_quality = "full" if events_path.exists() and events else "partial"856 return SnapshotData(857 repo=repo,858 snapshot_id=snapshot_id,859 snapshot_dir=snapshot_dir,860 manifest=manifest,861 issues=issues,862 pull_requests=pull_requests,863 comments=comments,864 reviews=reviews,865 review_comments=review_comments,866 pr_files=pr_files,867 pr_diffs=pr_diffs,868 links=links,869 events=events,870 evidence_quality=evidence_quality,871 )872 873 874async def _build_report(snapshot: SnapshotData, options: AnalysisOptions) -> AnalysisBuildResult:875 combined_links = _combined_links(snapshot)876 llm_available = _can_use_fast_agent()877 hybrid_review_cache = HybridReviewCacheStore(878 hybrid_review_cache_dir(snapshot.snapshot_dir),879 _hybrid_review_cache_manifest(),880 enabled=options.ranking_backend == "hybrid",881 )882 if hybrid_review_cache.invalidation_reason is not None:883 _analysis_log(884 "Hybrid review cache invalidated; ignoring cached entries "885 f"({hybrid_review_cache.invalidation_reason})"886 )887 issue_map = {int(row["number"]): row for row in snapshot.issues}888 pr_map = {int(row["number"]): row for row in snapshot.pull_requests}889 suppressed_pr_reasons = suppressed_pull_request_reasons(890 snapshot.pull_requests,891 snapshot.pr_files,892 compile_cluster_suppression_rules(options.cluster_suppression_rules),893 )894 if suppressed_pr_reasons:895 original_pr_count = len(pr_map)896 pr_map = {897 number: row for number, row in pr_map.items() if number not in suppressed_pr_reasons898 }899 _analysis_log(900 f"Suppressing {len(suppressed_pr_reasons)} routine PRs from clustering: "901 f"{len(pr_map)}/{original_pr_count} PRs kept"902 )903 if options.open_prs_only:904 original_pr_count = len(pr_map)905 pr_map = {906 number: row907 for number, row in pr_map.items()908 if str(row.get("state") or "").lower() == "open"909 }910 _analysis_log(911 f"Restricting PR analysis to open PRs only: {len(pr_map)}/{original_pr_count} PRs kept "912 "(draft PRs remain eligible)"913 )914 comment_map = {915 int(row["github_id"]): row for row in snapshot.comments if row.get("github_id") is not None916 }917 review_map = {918 int(row["github_id"]): row for row in snapshot.reviews if row.get("github_id") is not None919 }920 review_comment_map = {921 int(row["github_id"]): row922 for row in snapshot.review_comments923 if row.get("github_id") is not None924 }925 926 inbound_references, _ = _reference_counts(927 snapshot.repo,928 combined_links,929 issue_map,930 pr_map,931 comment_map,932 review_map,933 review_comment_map,934 )935 explicit_issue_link_targets = _explicit_pr_issue_targets(936 repo=snapshot.repo,937 combined_links=combined_links,938 issue_map=issue_map,939 pr_map=pr_map,940 )941 features = _artifact_features(942 snapshot,943 options=options,944 issue_map=issue_map,945 pr_map=pr_map,946 inbound_references=inbound_references,947 explicit_issue_link_targets=explicit_issue_link_targets,948 )949 issue_hard_pairs = _issue_hard_pairs(950 repo=snapshot.repo,951 combined_links=combined_links,952 issue_map=issue_map,953 pr_map=pr_map,954 comment_map=comment_map,955 review_map=review_map,956 review_comment_map=review_comment_map,957 )958 issue_soft_candidates = _issue_soft_candidates(issue_map, features, issue_hard_pairs)959 pr_soft_candidates, pr_pair_target_issues = _pr_duplicate_candidates(960 options=options,961 snapshot=snapshot,962 issue_map=issue_map,963 pr_map=pr_map,964 features=features,965 )966 review_semaphore = asyncio.Semaphore(options.hybrid_llm_concurrency)967 (968 (accepted_issue_pairs, issue_llm_enabled, issue_llm_reviews),969 (accepted_pr_pairs, pr_llm_enabled, pr_llm_reviews),970 ) = await asyncio.gather(971 _accepted_soft_pairs(972 options=options,973 snapshot=snapshot,974 features=features,975 hard_pairs=issue_hard_pairs,976 soft_candidates=issue_soft_candidates,977 label="issue",978 hybrid_review_cache=hybrid_review_cache,979 llm_available=llm_available,980 review_semaphore=review_semaphore,981 ),982 _accepted_soft_pairs(983 options=options,984 snapshot=snapshot,985 features=features,986 hard_pairs={},987 soft_candidates=pr_soft_candidates,988 label="pull_request",989 hybrid_review_cache=hybrid_review_cache,990 llm_available=llm_available,991 review_semaphore=review_semaphore,992 ),993 )994 issue_pairs = dict(issue_hard_pairs)995 for pair, detail in accepted_issue_pairs.items():996 issue_pairs.setdefault(pair, set()).update(997 detail.get("evidence_types") or {"soft_similarity"}998 )999 pr_pairs: dict[tuple[str, str], set[str]] = {}1000 for pair, detail in accepted_pr_pairs.items():1001 pr_pairs.setdefault(pair, set()).update(detail.get("evidence_types") or {"soft_similarity"})1002 1003 issue_clusters = _clusters(1004 snapshot=snapshot,1005 features=features,1006 final_pairs=issue_pairs,1007 pair_target_issues=defaultdict(set),1008 llm_cluster_payloads={},1009 )1010 pr_clusters = _clusters(1011 snapshot=snapshot,1012 features=features,1013 final_pairs=pr_pairs,1014 pair_target_issues=pr_pair_target_issues,1015 llm_cluster_payloads={},1016 )1017 clusters = _meta_bug_clusters(1018 features=features,1019 issue_clusters=issue_clusters,1020 pr_clusters=pr_clusters,1021 explicit_issue_link_targets=explicit_issue_link_targets,1022 issue_map=issue_map,1023 pr_map=pr_map,1024 )1025 1026 meta_clusters = sorted(1027 clusters, key=lambda cluster: (-cluster.cluster_score, cluster.cluster_id)1028 )[: options.max_clusters]1029 duplicate_issues = [1030 DuplicateIssuesEntry(1031 cluster_id=cluster.cluster_id,1032 canonical_issue_number=cluster.canonical_issue_number,1033 duplicate_issue_numbers=[1034 number1035 for number in cluster.issue_numbers1036 if number != cluster.canonical_issue_number1037 ],1038 reason=_duplicate_issue_reason(cluster),1039 )1040 for cluster in clusters1041 if cluster.canonical_issue_number is not None and len(cluster.issue_numbers) >= 21042 ]1043 duplicate_prs = [1044 DuplicatePrsEntry(1045 cluster_id=cluster.cluster_id,1046 canonical_pr_number=cluster.canonical_pr_number,1047 duplicate_pr_numbers=[1048 number for number in cluster.pr_numbers if number != cluster.canonical_pr_number1049 ],1050 target_issue_number=cluster.target_issue_number,1051 reason=_duplicate_pr_reason(cluster),1052 )1053 for cluster in clusters1054 if cluster.canonical_pr_number is not None and len(cluster.pr_numbers) >= 21055 ]1056 best_issue = _best_issue(meta_clusters, features)1057 best_pr = _best_pr(meta_clusters, features)1058 return AnalysisBuildResult(1059 report=AnalysisReport(1060 schema_version="1.0",1061 repo=snapshot.repo,1062 snapshot_id=snapshot.snapshot_id,1063 generated_at=_iso_now(),1064 evidence_quality=snapshot.evidence_quality,1065 llm_enrichment=issue_llm_enabled or pr_llm_enabled,1066 meta_bugs=[1067 MetaBugEntry(1068 cluster_id=cluster.cluster_id,1069 summary=cluster.summary,1070 status=cluster.status,1071 confidence=round(cluster.confidence, 3),1072 canonical_issue_number=cluster.canonical_issue_number,1073 canonical_pr_number=cluster.canonical_pr_number,1074 issue_numbers=cluster.issue_numbers,1075 pr_numbers=cluster.pr_numbers,1076 evidence_types=cluster.evidence_types,1077 pr_comparisons=_cluster_pr_comparisons(cluster, features),1078 )1079 for cluster in meta_clusters1080 ],1081 duplicate_issues=duplicate_issues,1082 duplicate_prs=duplicate_prs,1083 best_issue=best_issue,1084 best_pr=best_pr,1085 ),1086 llm_reviews=issue_llm_reviews + pr_llm_reviews,1087 )1088 1089 1090def _iso_now() -> str:1091 return datetime.now(tz=UTC).replace(microsecond=0).isoformat().replace("+00:00", "Z")1092 1093 1094def _combined_links(snapshot: SnapshotData) -> list[dict[str, Any]]:1095 owner, repo_name = snapshot.repo.split("/", 1)1096 extracted_at = snapshot.manifest.get("extracted_at") or _iso_now()1097 rows = list(snapshot.links)1098 for issue in snapshot.issues:1099 rows.extend(1100 build_text_link_rows(1101 repo=snapshot.repo,1102 owner=owner,1103 repo_name=repo_name,1104 source_type="issue",1105 source_number=int(issue["number"]),1106 source_id=issue.get("github_id"),1107 body=issue.get("body"),1108 snapshot_id=snapshot.snapshot_id,1109 extracted_at=extracted_at,1110 )1111 )1112 for pr in snapshot.pull_requests:1113 rows.extend(1114 build_text_link_rows(1115 repo=snapshot.repo,1116 owner=owner,1117 repo_name=repo_name,1118 source_type="pull_request",1119 source_number=int(pr["number"]),1120 source_id=pr.get("github_id"),1121 body=pr.get("body"),1122 snapshot_id=snapshot.snapshot_id,1123 extracted_at=extracted_at,1124 )1125 )1126 for comment in snapshot.comments:1127 if comment.get("parent_number") is None:1128 continue1129 rows.extend(1130 build_text_link_rows(1131 repo=snapshot.repo,1132 owner=owner,1133 repo_name=repo_name,1134 source_type="comment",1135 source_number=int(comment["parent_number"]),1136 source_id=comment.get("github_id"),1137 body=comment.get("body"),1138 snapshot_id=snapshot.snapshot_id,1139 extracted_at=extracted_at,1140 )1141 )1142 for review in snapshot.reviews:1143 rows.extend(1144 build_text_link_rows(1145 repo=snapshot.repo,1146 owner=owner,1147 repo_name=repo_name,1148 source_type="review",1149 source_number=int(review["pull_request_number"]),1150 source_id=review.get("github_id"),1151 body=review.get("body"),1152 snapshot_id=snapshot.snapshot_id,1153 extracted_at=extracted_at,1154 )1155 )1156 for review_comment in snapshot.review_comments:1157 rows.extend(1158 build_text_link_rows(1159 repo=snapshot.repo,1160 owner=owner,1161 repo_name=repo_name,1162 source_type="review_comment",1163 source_number=int(review_comment["pull_request_number"]),1164 source_id=review_comment.get("github_id"),1165 body=review_comment.get("body"),1166 snapshot_id=snapshot.snapshot_id,1167 extracted_at=extracted_at,1168 )1169 )1170 deduped: dict[tuple[Any, ...], dict[str, Any]] = {}1171 for row in rows:1172 key = tuple(row.get(field) for field in LINK_KEY_FIELDS)1173 deduped[key] = row1174 return list(deduped.values())1175 1176 1177def _reference_counts(1178 repo: str,1179 links: list[dict[str, Any]],1180 issue_map: dict[int, dict[str, Any]],1181 pr_map: dict[int, dict[str, Any]],1182 comment_map: dict[int, dict[str, Any]],1183 review_map: dict[int, dict[str, Any]],1184 review_comment_map: dict[int, dict[str, Any]],1185) -> tuple[Counter[str], defaultdict[int, set[int]]]:1186 inbound_references: Counter[str] = Counter()1187 explicit_issue_link_targets: defaultdict[int, set[int]] = defaultdict(set)1188 for row in links:1189 source_node = _resolve_source_node(1190 row, issue_map, pr_map, comment_map, review_map, review_comment_map1191 )1192 target_node = _resolve_target_node(repo, row, issue_map, pr_map)1193 if source_node is not None and target_node is not None:1194 inbound_references[target_node] += 11195 if (1196 source_node1197 and target_node1198 and source_node.startswith("pull_request:")1199 and target_node.startswith("issue:")1200 ):