evalstate/diffusers-pr-api
0
1from __future__ import annotations2 3import json4from pathlib import Path5from typing import Any6 7import duckdb8 9TABLE_COLUMNS: dict[str, tuple[str, ...]] = {10 "pr_search_runs": (11 "id",12 "repo",13 "snapshot_id",14 "snapshot_dir",15 "source_type",16 "hf_repo_id",17 "hf_revision",18 "started_at",19 "finished_at",20 "status",21 "settings_json",22 "notes",23 ),24 "pr_search_active_run": (25 "repo",26 "run_id",27 "activated_at",28 ),29 "pr_search_documents": (30 "run_id",31 "repo",32 "pr_number",33 "github_id",34 "author_login",35 "state",36 "draft",37 "merged",38 "title",39 "base_ref",40 "created_at",41 "updated_at",42 "merged_at",43 "additions",44 "deletions",45 "changed_files",46 "comments_count",47 "review_comments_count",48 "html_url",49 ),50 "pr_search_contributors": (51 "run_id",52 "repo",53 "snapshot_id",54 "report_generated_at",55 "window_days",56 "author_login",57 "name",58 "profile_url",59 "repo_pull_requests_url",60 "repo_issues_url",61 "repo_first_seen_at",62 "repo_last_seen_at",63 "repo_primary_artifact_count",64 "repo_artifact_count",65 "snapshot_issue_count",66 "snapshot_pr_count",67 "snapshot_comment_count",68 "snapshot_review_count",69 "snapshot_review_comment_count",70 "repo_association",71 "new_to_repo",72 "first_seen_in_snapshot",73 "report_reason",74 "account_age_days",75 "young_account",76 "follow_through_score",77 "breadth_score",78 "automation_risk_signal",79 "heuristic_note",80 "public_orgs_json",81 "visible_authored_pr_count",82 "merged_pr_count",83 "closed_unmerged_pr_count",84 "open_pr_count",85 "merged_pr_rate",86 "closed_unmerged_pr_rate",87 "still_open_pr_rate",88 "distinct_repos_with_authored_prs",89 "distinct_repos_with_open_prs",90 "fetch_error",91 ),92 "pr_scope_features": (93 "run_id",94 "repo",95 "pr_number",96 "feature_version",97 "total_changed_lines",98 "file_count",99 "directory_count",100 "dominant_dir_share",101 "filenames_json",102 "directories_json",103 "vector_json",104 "computed_at",105 ),106 "pr_scope_run_artifacts": (107 "run_id",108 "repo",109 "feature_version",110 "idf_json",111 "computed_at",112 ),113 "pr_scope_neighbors": (114 "run_id",115 "repo",116 "left_pr_number",117 "right_pr_number",118 "rank_from_left",119 "rank_from_right",120 "similarity",121 "content_similarity",122 "size_similarity",123 "breadth_similarity",124 "concentration_similarity",125 "shared_filenames_json",126 "shared_directories_json",127 "created_at",128 ),129 "pr_scope_clusters": (130 "run_id",131 "repo",132 "cluster_id",133 "representative_pr_number",134 "cluster_size",135 "average_similarity",136 "summary",137 "shared_filenames_json",138 "shared_directories_json",139 "created_at",140 ),141 "pr_scope_cluster_members": (142 "run_id",143 "repo",144 "cluster_id",145 "pr_number",146 "member_role",147 ),148 "pr_scope_cluster_candidates": (149 "run_id",150 "repo",151 "pr_number",152 "cluster_id",153 "candidate_rank",154 "candidate_score",155 "matched_member_count",156 "best_member_pr_number",157 "max_member_similarity",158 "avg_top_member_similarity",159 "evidence_json",160 "assigned",161 ),162}163 164 165SCHEMA_SQL = """166CREATE TABLE IF NOT EXISTS pr_search_runs (167 id VARCHAR,168 repo VARCHAR,169 snapshot_id VARCHAR,170 snapshot_dir VARCHAR,171 source_type VARCHAR,172 hf_repo_id VARCHAR,173 hf_revision VARCHAR,174 started_at VARCHAR,175 finished_at VARCHAR,176 status VARCHAR,177 settings_json VARCHAR,178 notes VARCHAR179);180CREATE TABLE IF NOT EXISTS pr_search_active_run (181 repo VARCHAR,182 run_id VARCHAR,183 activated_at VARCHAR184);185CREATE TABLE IF NOT EXISTS pr_search_documents (186 run_id VARCHAR,187 repo VARCHAR,188 pr_number BIGINT,189 github_id BIGINT,190 author_login VARCHAR,191 state VARCHAR,192 draft BOOLEAN,193 merged BOOLEAN,194 title VARCHAR,195 base_ref VARCHAR,196 created_at VARCHAR,197 updated_at VARCHAR,198 merged_at VARCHAR,199 additions BIGINT,200 deletions BIGINT,201 changed_files BIGINT,202 comments_count BIGINT,203 review_comments_count BIGINT,204 html_url VARCHAR205);206CREATE TABLE IF NOT EXISTS pr_search_contributors (207 run_id VARCHAR,208 repo VARCHAR,209 snapshot_id VARCHAR,210 report_generated_at VARCHAR,211 window_days BIGINT,212 author_login VARCHAR,213 name VARCHAR,214 profile_url VARCHAR,215 repo_pull_requests_url VARCHAR,216 repo_issues_url VARCHAR,217 repo_first_seen_at VARCHAR,218 repo_last_seen_at VARCHAR,219 repo_primary_artifact_count BIGINT,220 repo_artifact_count BIGINT,221 snapshot_issue_count BIGINT,222 snapshot_pr_count BIGINT,223 snapshot_comment_count BIGINT,224 snapshot_review_count BIGINT,225 snapshot_review_comment_count BIGINT,226 repo_association VARCHAR,227 new_to_repo BOOLEAN,228 first_seen_in_snapshot BOOLEAN,229 report_reason VARCHAR,230 account_age_days BIGINT,231 young_account BOOLEAN,232 follow_through_score VARCHAR,233 breadth_score VARCHAR,234 automation_risk_signal VARCHAR,235 heuristic_note VARCHAR,236 public_orgs_json VARCHAR,237 visible_authored_pr_count BIGINT,238 merged_pr_count BIGINT,239 closed_unmerged_pr_count BIGINT,240 open_pr_count BIGINT,241 merged_pr_rate DOUBLE,242 closed_unmerged_pr_rate DOUBLE,243 still_open_pr_rate DOUBLE,244 distinct_repos_with_authored_prs BIGINT,245 distinct_repos_with_open_prs BIGINT,246 fetch_error VARCHAR247);248CREATE TABLE IF NOT EXISTS pr_scope_features (249 run_id VARCHAR,250 repo VARCHAR,251 pr_number BIGINT,252 feature_version VARCHAR,253 total_changed_lines BIGINT,254 file_count BIGINT,255 directory_count BIGINT,256 dominant_dir_share DOUBLE,257 filenames_json VARCHAR,258 directories_json VARCHAR,259 vector_json VARCHAR,260 computed_at VARCHAR261);262CREATE TABLE IF NOT EXISTS pr_scope_run_artifacts (263 run_id VARCHAR,264 repo VARCHAR,265 feature_version VARCHAR,266 idf_json VARCHAR,267 computed_at VARCHAR268);269CREATE TABLE IF NOT EXISTS pr_scope_neighbors (270 run_id VARCHAR,271 repo VARCHAR,272 left_pr_number BIGINT,273 right_pr_number BIGINT,274 rank_from_left BIGINT,275 rank_from_right BIGINT,276 similarity DOUBLE,277 content_similarity DOUBLE,278 size_similarity DOUBLE,279 breadth_similarity DOUBLE,280 concentration_similarity DOUBLE,281 shared_filenames_json VARCHAR,282 shared_directories_json VARCHAR,283 created_at VARCHAR284);285CREATE TABLE IF NOT EXISTS pr_scope_clusters (286 run_id VARCHAR,287 repo VARCHAR,288 cluster_id VARCHAR,289 representative_pr_number BIGINT,290 cluster_size BIGINT,291 average_similarity DOUBLE,292 summary VARCHAR,293 shared_filenames_json VARCHAR,294 shared_directories_json VARCHAR,295 created_at VARCHAR296);297CREATE TABLE IF NOT EXISTS pr_scope_cluster_members (298 run_id VARCHAR,299 repo VARCHAR,300 cluster_id VARCHAR,301 pr_number BIGINT,302 member_role VARCHAR303);304CREATE TABLE IF NOT EXISTS pr_scope_cluster_candidates (305 run_id VARCHAR,306 repo VARCHAR,307 pr_number BIGINT,308 cluster_id VARCHAR,309 candidate_rank BIGINT,310 candidate_score DOUBLE,311 matched_member_count BIGINT,312 best_member_pr_number BIGINT,313 max_member_similarity DOUBLE,314 avg_top_member_similarity DOUBLE,315 evidence_json VARCHAR,316 assigned BOOLEAN317);318CREATE INDEX IF NOT EXISTS idx_pr_search_active_run_repo ON pr_search_active_run (repo);319CREATE INDEX IF NOT EXISTS idx_pr_search_runs_repo_status ON pr_search_runs (repo, status);320CREATE INDEX IF NOT EXISTS idx_pr_search_documents_run_pr ON pr_search_documents (run_id, pr_number);321CREATE INDEX IF NOT EXISTS idx_pr_search_documents_run_author ON pr_search_documents (run_id, author_login);322CREATE INDEX IF NOT EXISTS idx_pr_search_contributors_run_author ON pr_search_contributors (run_id, author_login);323CREATE INDEX IF NOT EXISTS idx_pr_scope_features_run_pr ON pr_scope_features (run_id, pr_number);324CREATE INDEX IF NOT EXISTS idx_pr_scope_run_artifacts_run ON pr_scope_run_artifacts (run_id);325CREATE INDEX IF NOT EXISTS idx_pr_scope_neighbors_run_left ON pr_scope_neighbors (run_id, left_pr_number);326CREATE INDEX IF NOT EXISTS idx_pr_scope_neighbors_run_right ON pr_scope_neighbors (run_id, right_pr_number);327CREATE INDEX IF NOT EXISTS idx_pr_scope_clusters_run_cluster ON pr_scope_clusters (run_id, cluster_id);328CREATE INDEX IF NOT EXISTS idx_pr_scope_cluster_members_run_pr ON pr_scope_cluster_members (run_id, pr_number);329CREATE INDEX IF NOT EXISTS idx_pr_scope_cluster_candidates_run_pr ON pr_scope_cluster_candidates (run_id, pr_number);330"""331 332 333def connect_pr_search_db(path: Path, *, read_only: bool = False) -> duckdb.DuckDBPyConnection:334 resolved = path.resolve()335 if read_only and not resolved.exists():336 raise FileNotFoundError(f"PR search database does not exist: {resolved}")337 if not read_only:338 resolved.parent.mkdir(parents=True, exist_ok=True)339 connection = duckdb.connect(str(resolved), read_only=read_only)340 if not read_only:341 ensure_pr_search_schema(connection)342 return connection343 344 345def ensure_pr_search_schema(connection: duckdb.DuckDBPyConnection) -> None:346 connection.execute(SCHEMA_SQL)347 connection.execute(348 "ALTER TABLE pr_search_documents ADD COLUMN IF NOT EXISTS author_login VARCHAR"349 )350 351 352def insert_rows(353 connection: duckdb.DuckDBPyConnection,354 table_name: str,355 rows: list[dict[str, Any]],356) -> None:357 if not rows:358 return359 columns = TABLE_COLUMNS[table_name]360 placeholders = ", ".join("?" for _ in columns)361 column_sql = ", ".join(columns)362 values = [tuple(_db_value(row.get(column)) for column in columns) for row in rows]363 connection.executemany(364 f"INSERT INTO {table_name} ({column_sql}) VALUES ({placeholders})",365 values,366 )367 368 369def update_run_status(370 connection: duckdb.DuckDBPyConnection,371 *,372 run_id: str,373 status: str,374 finished_at: str | None = None,375 notes: str | None = None,376) -> None:377 connection.execute(378 """379 UPDATE pr_search_runs380 SET status = ?, finished_at = COALESCE(?, finished_at), notes = COALESCE(?, notes)381 WHERE id = ?382 """,383 [status, finished_at, notes, run_id],384 )385 386 387def replace_active_run(388 connection: duckdb.DuckDBPyConnection,389 *,390 repo: str,391 run_id: str,392 activated_at: str,393) -> str | None:394 previous = fetch_one(395 connection,396 "SELECT run_id FROM pr_search_active_run WHERE repo = ?",397 [repo],398 )399 connection.execute("DELETE FROM pr_search_active_run WHERE repo = ?", [repo])400 connection.execute(401 "INSERT INTO pr_search_active_run (repo, run_id, activated_at) VALUES (?, ?, ?)",402 [repo, run_id, activated_at],403 )404 previous_run_id = None if previous is None else str(previous["run_id"])405 if previous_run_id is not None and previous_run_id != run_id:406 connection.execute(407 "UPDATE pr_search_runs SET status = 'superseded' WHERE id = ? AND status = 'succeeded'",408 [previous_run_id],409 )410 return previous_run_id411 412 413def resolve_active_run(414 connection: duckdb.DuckDBPyConnection,415 *,416 repo: str | None = None,417) -> dict[str, Any]:418 if repo is None:419 active_repos = fetch_rows(420 connection,421 "SELECT repo FROM pr_search_active_run ORDER BY repo",422 )423 if not active_repos:424 raise ValueError("No active PR search run found.")425 if len(active_repos) > 1:426 raise ValueError("Multiple active repos found; pass --repo.")427 repo = str(active_repos[0]["repo"])428 row = fetch_one(429 connection,430 """431 SELECT r.*432 FROM pr_search_runs AS r433 JOIN pr_search_active_run AS a434 ON a.run_id = r.id AND a.repo = r.repo435 WHERE a.repo = ?436 """,437 [repo],438 )439 if row is None:440 raise ValueError(f"No active PR search run found for repo {repo!r}.")441 return row442 443 444def get_run_counts(connection: duckdb.DuckDBPyConnection, *, run_id: str) -> dict[str, int]:445 return {446 "documents": _count(connection, "pr_search_documents", run_id),447 "contributors": _count(connection, "pr_search_contributors", run_id),448 "features": _count(connection, "pr_scope_features", run_id),449 "run_artifacts": _count(connection, "pr_scope_run_artifacts", run_id),450 "neighbors": _count(connection, "pr_scope_neighbors", run_id),451 "clusters": _count(connection, "pr_scope_clusters", run_id),452 "cluster_members": _count(connection, "pr_scope_cluster_members", run_id),453 "cluster_candidates": _count(connection, "pr_scope_cluster_candidates", run_id),454 }455 456 457def get_document(458 connection: duckdb.DuckDBPyConnection,459 *,460 run_id: str,461 pr_number: int,462) -> dict[str, Any] | None:463 return fetch_one(464 connection,465 "SELECT * FROM pr_search_documents WHERE run_id = ? AND pr_number = ?",466 [run_id, pr_number],467 )468 469 470def get_contributor(471 connection: duckdb.DuckDBPyConnection,472 *,473 run_id: str,474 author_login: str,475) -> dict[str, Any] | None:476 return fetch_one(477 connection,478 """479 SELECT *480 FROM pr_search_contributors481 WHERE run_id = ? AND lower(author_login) = lower(?)482 """,483 [run_id, author_login],484 )485 486 487def get_contributor_pulls(488 connection: duckdb.DuckDBPyConnection,489 *,490 run_id: str,491 author_login: str,492 limit: int,493) -> list[dict[str, Any]]:494 return fetch_rows(495 connection,496 """497 SELECT498 pr_number,499 github_id,500 author_login,501 state,502 draft,503 merged,504 title,505 base_ref,506 created_at,507 updated_at,508 merged_at,509 additions,510 deletions,511 changed_files,512 comments_count,513 review_comments_count,514 html_url515 FROM pr_search_documents516 WHERE run_id = ? AND lower(author_login) = lower(?)517 ORDER BY updated_at DESC NULLS LAST, pr_number DESC518 LIMIT ?519 """,520 [run_id, author_login, limit],521 )522 523 524def get_feature(525 connection: duckdb.DuckDBPyConnection,526 *,527 run_id: str,528 pr_number: int,529) -> dict[str, Any] | None:530 return fetch_one(531 connection,532 "SELECT * FROM pr_scope_features WHERE run_id = ? AND pr_number = ?",533 [run_id, pr_number],534 )535 536 537def get_scope_run_artifact(538 connection: duckdb.DuckDBPyConnection,539 *,540 run_id: str,541) -> dict[str, Any] | None:542 try:543 return fetch_one(544 connection,545 """546 SELECT *547 FROM pr_scope_run_artifacts548 WHERE run_id = ?549 """,550 [run_id],551 )552 except duckdb.Error:553 return None554 555 556def get_similar_pr_rows(557 connection: duckdb.DuckDBPyConnection,558 *,559 run_id: str,560 pr_number: int,561 limit: int,562) -> list[dict[str, Any]]:563 return fetch_rows(564 connection,565 """566 SELECT567 CASE WHEN left_pr_number = ? THEN right_pr_number ELSE left_pr_number END AS neighbor_pr_number,568 CASE WHEN left_pr_number = ? THEN rank_from_left ELSE rank_from_right END AS neighbor_rank,569 similarity,570 content_similarity,571 size_similarity,572 breadth_similarity,573 concentration_similarity,574 shared_filenames_json,575 shared_directories_json576 FROM pr_scope_neighbors577 WHERE run_id = ? AND (? = left_pr_number OR ? = right_pr_number)578 ORDER BY neighbor_rank IS NULL, neighbor_rank, similarity DESC, neighbor_pr_number579 LIMIT ?580 """,581 [pr_number, pr_number, run_id, pr_number, pr_number, limit],582 )583 584 585def get_candidate_cluster_rows(586 connection: duckdb.DuckDBPyConnection,587 *,588 run_id: str,589 pr_number: int,590 limit: int,591) -> list[dict[str, Any]]:592 return fetch_rows(593 connection,594 """595 SELECT596 c.cluster_id,597 c.candidate_rank,598 c.candidate_score,599 c.matched_member_count,600 c.best_member_pr_number,601 c.max_member_similarity,602 c.avg_top_member_similarity,603 c.evidence_json,604 c.assigned,605 cl.representative_pr_number,606 cl.cluster_size,607 cl.average_similarity,608 cl.summary,609 cl.shared_filenames_json,610 cl.shared_directories_json,611 d.title AS representative_title612 FROM pr_scope_cluster_candidates AS c613 JOIN pr_scope_clusters AS cl614 ON cl.run_id = c.run_id AND cl.cluster_id = c.cluster_id615 LEFT JOIN pr_search_documents AS d616 ON d.run_id = cl.run_id AND d.pr_number = cl.representative_pr_number617 WHERE c.run_id = ? AND c.pr_number = ?618 ORDER BY c.candidate_rank, c.candidate_score DESC, c.cluster_id619 LIMIT ?620 """,621 [run_id, pr_number, limit],622 )623 624 625def get_cluster(626 connection: duckdb.DuckDBPyConnection,627 *,628 run_id: str,629 cluster_id: str,630) -> dict[str, Any] | None:631 return fetch_one(632 connection,633 "SELECT * FROM pr_scope_clusters WHERE run_id = ? AND cluster_id = ?",634 [run_id, cluster_id],635 )636 637 638def get_cluster_members(639 connection: duckdb.DuckDBPyConnection,640 *,641 run_id: str,642 cluster_id: str,643) -> list[dict[str, Any]]:644 return fetch_rows(645 connection,646 """647 SELECT648 m.pr_number,649 m.member_role,650 d.title,651 d.html_url,652 d.state,653 d.draft654 FROM pr_scope_cluster_members AS m655 LEFT JOIN pr_search_documents AS d656 ON d.run_id = m.run_id AND d.pr_number = m.pr_number657 WHERE m.run_id = ? AND m.cluster_id = ?658 ORDER BY m.member_role != 'representative', m.pr_number659 """,660 [run_id, cluster_id],661 )662 663 664def get_cluster_ids_for_prs(665 connection: duckdb.DuckDBPyConnection,666 *,667 run_id: str,668 pr_numbers: list[int],669) -> dict[int, list[str]]:670 if not pr_numbers:671 return {}672 placeholders = ", ".join("?" for _ in pr_numbers)673 rows = fetch_rows(674 connection,675 f"""676 SELECT pr_number, cluster_id677 FROM pr_scope_cluster_members678 WHERE run_id = ? AND pr_number IN ({placeholders})679 ORDER BY pr_number, cluster_id680 """,681 [run_id, *pr_numbers],682 )683 result: dict[int, list[str]] = {}684 for row in rows:685 result.setdefault(int(row["pr_number"]), []).append(str(row["cluster_id"]))686 return result687 688 689def get_shared_cluster_ids(690 connection: duckdb.DuckDBPyConnection,691 *,692 run_id: str,693 left_pr_number: int,694 right_pr_number: int,695) -> list[str]:696 rows = fetch_rows(697 connection,698 """699 SELECT left_members.cluster_id700 FROM pr_scope_cluster_members AS left_members701 JOIN pr_scope_cluster_members AS right_members702 ON right_members.run_id = left_members.run_id703 AND right_members.cluster_id = left_members.cluster_id704 WHERE left_members.run_id = ?705 AND left_members.pr_number = ?706 AND right_members.pr_number = ?707 ORDER BY left_members.cluster_id708 """,709 [run_id, left_pr_number, right_pr_number],710 )711 return [str(row["cluster_id"]) for row in rows]712 713 714def get_pair_neighbor_row(715 connection: duckdb.DuckDBPyConnection,716 *,717 run_id: str,718 left_pr_number: int,719 right_pr_number: int,720) -> dict[str, Any] | None:721 canonical_left = min(left_pr_number, right_pr_number)722 canonical_right = max(left_pr_number, right_pr_number)723 return fetch_one(724 connection,725 """726 SELECT *727 FROM pr_scope_neighbors728 WHERE run_id = ? AND left_pr_number = ? AND right_pr_number = ?729 """,730 [run_id, canonical_left, canonical_right],731 )732 733 734def fetch_rows(735 connection: duckdb.DuckDBPyConnection,736 sql: str,737 parameters: list[Any] | tuple[Any, ...] | None = None,738) -> list[dict[str, Any]]:739 cursor = connection.execute(sql, parameters or [])740 columns = [column[0] for column in cursor.description]741 return [dict(zip(columns, row, strict=False)) for row in cursor.fetchall()]742 743 744def fetch_one(745 connection: duckdb.DuckDBPyConnection,746 sql: str,747 parameters: list[Any] | tuple[Any, ...] | None = None,748) -> dict[str, Any] | None:749 rows = fetch_rows(connection, sql, parameters)750 return rows[0] if rows else None751 752 753def _count(connection: duckdb.DuckDBPyConnection, table_name: str, run_id: str) -> int:754 row = fetch_one(755 connection,756 f"SELECT COUNT(*) AS row_count FROM {table_name} WHERE run_id = ?",757 [run_id],758 )759 return 0 if row is None else int(row["row_count"])760 761 762def _db_value(value: Any) -> Any:763 if isinstance(value, (dict, list)):764 return json.dumps(value, sort_keys=True)765 return value766 