from __future__ import annotations from collections import Counter, defaultdict from datetime import datetime from pathlib import Path from typing import Any from common import append_jsonl, ensure_utc, parse_datetime, read_jsonl, write_jsonl WAITING_STATUSES = {"waiting", "queued"} RUNNING_STATUSES = {"running", "processing"} HISTORY_STATS_THRESHOLD_DEFAULT = 500 HISTORY_ARCHIVE_LIMIT_DEFAULT = 5000 def load_ledger(path: Path) -> list[dict[str, Any]]: return read_jsonl(path) def append_ledger_entry(path: Path, entry: dict[str, Any]) -> None: append_jsonl(path, entry) def load_history_archive(path: Path) -> list[dict[str, Any]]: return read_jsonl(path) def update_history_archive( path: Path, tasks: list[dict[str, Any]], *, limit: int = HISTORY_ARCHIVE_LIMIT_DEFAULT, ) -> list[dict[str, Any]]: merged = merge_history_records(load_history_archive(path), tasks, limit=limit) path.parent.mkdir(parents=True, exist_ok=True) write_jsonl(path, merged) return merged def merge_history_records( existing: list[dict[str, Any]], new_records: list[dict[str, Any]], *, limit: int = HISTORY_ARCHIVE_LIMIT_DEFAULT, ) -> list[dict[str, Any]]: by_key: dict[str, dict[str, Any]] = {} for record in [*existing, *new_records]: key = history_record_key(record) previous = by_key.get(key) if previous is None or history_record_sort_key(record) >= history_record_sort_key(previous): by_key[key] = record merged = sorted(by_key.values(), key=history_record_sort_key, reverse=True) if limit > 0: merged = merged[:limit] return merged def history_record_key(record: dict[str, Any]) -> str: task_id = record.get("taskId") if task_id is not None: return f"task:{task_id}" model_id = record.get("modelId") or record.get("modelAddress") or "unknown" gpu_type = record.get("gpuType") or record.get("targetGpu") or "unknown" framework = record.get("framework") or "unknown" return f"fallback:{model_id}|{gpu_type}|{framework}" def history_record_sort_key(record: dict[str, Any]) -> tuple[float, str]: timestamp = ( parse_datetime(record.get("updateTime")) or parse_datetime(record.get("createTime")) or parse_datetime(record.get("submitTime")) ) return ((timestamp.timestamp() if timestamp else 0.0), str(record.get("taskId") or "")) def should_use_history_stats(ledger_entries: list[dict[str, Any]], threshold: int = HISTORY_STATS_THRESHOLD_DEFAULT) -> bool: return len(ledger_entries) >= threshold def build_empty_pre_submit_report( *, window_days: int, generated_at: datetime, task_type: str | None = None, target_gpu: str | None = None, reason: str | None = None, ledger_entries: int = 0, ) -> dict[str, Any]: empty_summary = summarize_records([]) return { "generatedAt": generated_at.isoformat(), "statsWindowDays": window_days, "taskType": task_type, "taskTypeSummary": empty_summary, "taskTypeSummaries": {}, "targetGpu": target_gpu, "targetGpuSummary": empty_summary, "gpuSummaries": {}, "totals": empty_summary, "combinationStats": {}, "warnings": [] if reason is None else [reason], "recordsAnalyzed": 0, "ledgerEntries": ledger_entries, "historyStatsEnabled": False, "historyStatsReason": reason, } def build_pre_submit_report( *, tasks: list[dict[str, Any]], ledger_entries: list[dict[str, Any]], window_days: int, generated_at: datetime, target_gpu: str | None = None, task_type: str | None = None, ) -> dict[str, Any]: ledger_by_task_id = {str(entry["taskId"]): entry for entry in ledger_entries if entry.get("taskId")} enriched_tasks = [enrich_task(task, ledger_by_task_id) for task in tasks] grouped_by_combo: dict[tuple[str, str, str], list[dict[str, Any]]] = defaultdict(list) grouped_by_gpu: dict[str, list[dict[str, Any]]] = defaultdict(list) grouped_by_task_type: dict[str, list[dict[str, Any]]] = defaultdict(list) for task in enriched_tasks: combo_key = ( task.get("taskType") or "unknown", task.get("gpuType") or "unknown", task.get("framework") or "unknown", ) grouped_by_combo[combo_key].append(task) grouped_by_gpu[combo_key[1]].append(task) grouped_by_task_type[combo_key[0]].append(task) combination_stats = { "|".join(combo_key): { "taskType": combo_key[0], "gpuType": combo_key[1], "framework": combo_key[2], **summarize_records(records), } for combo_key, records in grouped_by_combo.items() } gpu_summaries = { gpu_type: summarize_records(records) for gpu_type, records in grouped_by_gpu.items() } task_type_summaries = { current_task_type: summarize_records(records) for current_task_type, records in grouped_by_task_type.items() } warnings = build_warnings(gpu_summaries, combination_stats) selected_gpu_summary = gpu_summaries.get(target_gpu, summarize_records([])) if target_gpu else None selected_task_summary = task_type_summaries.get(task_type, summarize_records([])) if task_type else None return { "generatedAt": generated_at.isoformat(), "statsWindowDays": window_days, "taskType": task_type, "taskTypeSummary": selected_task_summary, "taskTypeSummaries": task_type_summaries, "targetGpu": target_gpu, "targetGpuSummary": selected_gpu_summary, "gpuSummaries": gpu_summaries, "totals": summarize_records(enriched_tasks), "combinationStats": combination_stats, "warnings": warnings, "recordsAnalyzed": len(enriched_tasks), "ledgerEntries": len(ledger_entries), } def enrich_task(task: dict[str, Any], ledger_by_task_id: dict[str, dict[str, Any]]) -> dict[str, Any]: enriched = dict(task) task_id = enriched.get("taskId") ledger_entry = ledger_by_task_id.get(str(task_id)) if task_id is not None else None if ledger_entry: enriched.setdefault("taskType", ledger_entry.get("taskType")) enriched.setdefault("framework", ledger_entry.get("framework")) enriched.setdefault("templateId", ledger_entry.get("templateId")) enriched.setdefault("taskType", "unknown") enriched.setdefault("framework", "unknown") return enriched def summarize_records(records: list[dict[str, Any]]) -> dict[str, Any]: total = len(records) success_count = 0 failure_count = 0 pending_count = 0 running_count = 0 failure_breakdown: Counter[str] = Counter() status_counter: Counter[str] = Counter() verify_counter: Counter[str] = Counter() for record in records: status = str(record.get("status") or "unknown").lower() verify_result = record.get("verifyResult") status_counter[status] += 1 verify_counter[str(verify_result)] += 1 if is_success(record): success_count += 1 elif is_failure(record): failure_count += 1 failure_breakdown[classify_failure(record)] += 1 else: pending_count += 1 failure_breakdown[classify_failure(record)] += 1 if status in RUNNING_STATUSES: running_count += 1 waiting_count = sum(status_counter[status] for status in WAITING_STATUSES) return { "total": total, "successCount": success_count, "failureCount": failure_count, "pendingCount": pending_count, "runningCount": running_count, "waitingCount": waiting_count, "successRate": ratio(success_count, total), "failureRate": ratio(failure_count, total), "pendingRate": ratio(pending_count, total), "statusCounts": dict(status_counter), "verifyResultCounts": dict(verify_counter), "failureBreakdown": dict(failure_breakdown), } def is_success(record: dict[str, Any]) -> bool: verify_result = record.get("verifyResult") status = str(record.get("status") or "").lower() return verify_result is not None and verify_result > 0 and status == "success" def is_failure(record: dict[str, Any]) -> bool: verify_result = record.get("verifyResult") status = str(record.get("status") or "").lower() if verify_result is not None and verify_result < 0: return True return status in {"failed", "error", "cancelled"} def classify_failure(record: dict[str, Any]) -> str: status = str(record.get("status") or "").lower() verify_result = record.get("verifyResult") if status in {"failed", "error", "cancelled"}: return "参数/模板问题" if status in WAITING_STATUSES or status in RUNNING_STATUSES or verify_result is None: return "排队中/未知" if verify_result is not None and verify_result < 0: if not record.get("logCosUrl") and not record.get("logSyncStatus"): return "日志缺失" return "验证失败" return "排队中/未知" def build_warnings(gpu_summaries: dict[str, dict[str, Any]], combination_stats: dict[str, dict[str, Any]]) -> list[str]: warnings: list[str] = [] for gpu_type, target_gpu_summary in gpu_summaries.items(): if target_gpu_summary["total"] >= 4 and target_gpu_summary["failureRate"] >= 0.5: warnings.append(f"GPU {gpu_type} 最近 7 天失败率偏高。") if target_gpu_summary["total"] >= 4 and target_gpu_summary["pendingRate"] >= 0.5: warnings.append(f"GPU {gpu_type} 最近 7 天排队/未知任务占比较高。") for stat in combination_stats.values(): if stat["framework"] == "unknown": continue if stat["total"] >= 3 and stat["failureRate"] >= 0.6: warnings.append(f"框架 {stat['framework']} 在 {stat['gpuType']} 的任务 {stat['taskType']} 上近期失败集中。") deduped: list[str] = [] for warning in warnings: if warning not in deduped: deduped.append(warning) return deduped def score_candidate(report: dict[str, Any], *, framework: str, task_type: str, target_gpu: str) -> tuple[float, list[str]]: key = f"{task_type}|{target_gpu}|{framework}" combo = (report.get("combinationStats") or {}).get(key) gpu_summaries = report.get("gpuSummaries") or {} target_gpu_summary = gpu_summaries.get(target_gpu) or {} task_type_summary = (report.get("taskTypeSummaries") or {}).get(task_type) or {} score = 100.0 warnings = list(report.get("warnings") or []) if combo: score += combo["successRate"] * 25.0 score -= combo["failureRate"] * 40.0 score -= combo["pendingRate"] * 15.0 if target_gpu_summary: score -= target_gpu_summary.get("pendingRate", 0.0) * 10.0 score -= target_gpu_summary.get("failureRate", 0.0) * 10.0 if task_type_summary: score += task_type_summary.get("successRate", 0.0) * 10.0 score -= task_type_summary.get("failureRate", 0.0) * 10.0 return round(score, 2), warnings def count_submissions_for_day( *, tasks: list[dict[str, Any]], ledger_entries: list[dict[str, Any]], day_start: datetime, day_end: datetime, ) -> dict[str, Any]: day_start = ensure_utc(day_start) day_end = ensure_utc(day_end) task_ids = { str(task.get("taskId")) for task in tasks if task.get("taskId") is not None } task_count = len(tasks) ledger_only = 0 for entry in ledger_entries: submit_time = parse_datetime(entry.get("submitTime")) if submit_time is None or submit_time < day_start or submit_time > day_end: continue task_id = str(entry.get("taskId")) if entry.get("taskId") is not None else None if task_id and task_id in task_ids: continue ledger_only += 1 return { "platformTaskCount": task_count, "ledgerOnlyCount": ledger_only, "totalCount": task_count + ledger_only, } def ratio(part: int, whole: int) -> float: if whole <= 0: return 0.0 return round(part / whole, 4)