From cc41ff412748347782bdaa1c52a357ee087aecb7 Mon Sep 17 00:00:00 2001 From: ethanwee1 Date: Wed, 16 Sep 2026 19:01:10 +0000 Subject: [PATCH 1/3] [CI] Resolve MI350 parity from canonical trunk Use trunk as the sole MI350 source and derive each run's actual CUDA and ROCm job families so current and baseline commits survive job-name and shard-count changes. --- .../download_testlogs | 274 +++++++++++++----- .../job_name_match.py | 58 ++++ .../parity_job_config.json | 14 +- .../test_job_name_match.py | 98 +++++++ 4 files changed, 354 insertions(+), 90 deletions(-) create mode 100644 .automation_scripts/pytorch-unit-test-scripts/job_name_match.py create mode 100644 .automation_scripts/pytorch-unit-test-scripts/test_job_name_match.py diff --git a/.automation_scripts/pytorch-unit-test-scripts/download_testlogs b/.automation_scripts/pytorch-unit-test-scripts/download_testlogs index 84dcb932625f2..c21da09fdd33a 100755 --- a/.automation_scripts/pytorch-unit-test-scripts/download_testlogs +++ b/.automation_scripts/pytorch-unit-test-scripts/download_testlogs @@ -9,6 +9,7 @@ try: import requests import re import sys + from job_name_match import choose_test_job_family from upload_stats_lib import unzip from upload_test_stats import download_gha_artifacts, download_s3_artifacts except ImportError: @@ -233,6 +234,22 @@ def get_cuda_test_jobs(jobs, cuda_job_prefix): return test_kind, test_jobs return CUDA_TEST_KINDS[0], [] +def resolve_test_job_family(wf, test_config, platform, configured_prefix, + fallback_shards): + """Resolve a run's actual job prefix and shard total.""" + jobs = get_workflow_jobs(wf) + family = choose_test_job_family( + jobs, test_config, platform, configured_prefix + ) + if family is None: + return configured_prefix, fallback_shards, jobs + if family["prefix"] != configured_prefix: + print( + f"NOTE: configured {platform} prefix '{configured_prefix}' is absent; " + f"using '{family['prefix']}' from run {wf['id']}" + ) + return family["prefix"], family["total"], jobs + def get_cuda_inductor_test_jobs(jobs): """Return the CUDA inductor test kind and jobs for either main or PR CI layouts.""" inductor_prefix = CUDA_JOB_PREFIXES["inductor"] @@ -546,7 +563,8 @@ def resolve_full_trunk_run(sha, cuda_job_prefix, status='success', ignore_status full push runs the CUDA test jobs, so anchor on the run that contains them (falling back to the newest completed run when none can be found). - Returns (trunk_wf, cuda_test_jobs, cuda_test_kind, all_cuda_jobs). + Returns (trunk_wf, cuda_test_jobs, cuda_test_kind, all_cuda_jobs, + resolved_cuda_job_prefix). """ params = {'per_page': 10, 'head_sha': sha} if not ignore_status and status: @@ -570,12 +588,15 @@ def resolve_full_trunk_run(sha, cuda_job_prefix, status='success', ignore_status for run in trunk_runs: jobs = get_workflow_jobs(run) - kind, test_jobs = get_cuda_test_jobs(jobs, cuda_job_prefix) - if test_jobs: + family = choose_test_job_family( + jobs, "default", "cuda", cuda_job_prefix + ) + if family: trunk_wf = run all_cuda_jobs = jobs - cuda_test_jobs = test_jobs - cuda_test_kind = kind + cuda_test_jobs = family["jobs"] + cuda_test_kind = family["kind"] + cuda_job_prefix = family["prefix"] print(f"Found CUDA test jobs in trunk run {run['id']}") break @@ -583,9 +604,14 @@ def resolve_full_trunk_run(sha, cuda_job_prefix, status='success', ignore_status # CUDA jobs may live in a run the jobs API does not surface (e.g. a # reusable workflow); the check-runs API resolves the actual run. print("No CUDA test jobs in any trunk run's jobs API, trying check-runs API...") - check_runs = get_check_runs_for_commit(sha, cuda_job_prefix) - cuda_test_kind, cuda_test_jobs = get_cuda_test_jobs(check_runs, cuda_job_prefix) - if cuda_test_jobs: + check_runs = get_check_runs_for_commit(sha, "") + family = choose_test_job_family( + check_runs, "default", "cuda", cuda_job_prefix + ) + if family: + cuda_test_kind = family["kind"] + cuda_test_jobs = family["jobs"] + cuda_job_prefix = family["prefix"] run_match = re.search(r'/runs/(\d+)/', cuda_test_jobs[0].get('details_url', '')) if run_match: actual_run_id = int(run_match.group(1)) @@ -602,7 +628,13 @@ def resolve_full_trunk_run(sha, cuda_job_prefix, status='success', ignore_status if trunk_wf is None and trunk_runs: trunk_wf = trunk_runs[0] - return trunk_wf, cuda_test_jobs, cuda_test_kind, all_cuda_jobs + return ( + trunk_wf, + cuda_test_jobs, + cuda_test_kind, + all_cuda_jobs, + cuda_job_prefix, + ) def main(): @@ -691,11 +723,13 @@ def main(): trunk_cuda_test_jobs = [] trunk_cuda_test_kind = CUDA_TEST_KINDS[0] trunk_all_cuda_jobs = [] + trunk_cuda_job_prefix = CUDA_JOB_PREFIXES["default"] if (not args.no_cuda) or arch_uses_trunk or fallback_uses_trunk: print("==============================================") print(f"Resolving canonical trunk run for sha: {sha}") print("==============================================") - trunk_full_wf, trunk_cuda_test_jobs, trunk_cuda_test_kind, trunk_all_cuda_jobs = \ + trunk_full_wf, trunk_cuda_test_jobs, trunk_cuda_test_kind, \ + trunk_all_cuda_jobs, trunk_cuda_job_prefix = \ resolve_full_trunk_run(sha, CUDA_JOB_PREFIXES["default"], status, args.ignore_status) if trunk_full_wf is not None: print(f"Using trunk run {trunk_full_wf['id']} as the canonical run for this SHA " @@ -710,10 +744,13 @@ def main(): error_msg="Error: Periodic workflow not found in scanned workflow runs." #https://docs.github.com/en/rest/actions/workflow-runs#list-workflow-runs-for-a-repository periodic_fallback_used = False - try: - periodic_wf = download_workflow_run(created=args.created, max_pages=args.max_pages, workflow=ROCmWorkflowNames["distributed"], sha=periodic_sha, ignore_status=args.ignore_status, status=status, error_msg=error_msg) - except (IndexError, Exception): - periodic_wf = None + if ROCmWorkflowNames["distributed"] == "trunk" and trunk_full_wf is not None: + periodic_wf = trunk_full_wf + else: + try: + periodic_wf = download_workflow_run(created=args.created, max_pages=args.max_pages, workflow=ROCmWorkflowNames["distributed"], sha=periodic_sha, ignore_status=args.ignore_status, status=status, error_msg=error_msg) + except (IndexError, Exception): + periodic_wf = None periodic_fallbacks = _fallbacks_for("distributed") if periodic_wf is None and arch in periodic_fallbacks: fallback_wf, fallback_prefix = periodic_fallbacks[arch] @@ -735,6 +772,13 @@ def main(): else: dist_job_prefix = rocm_job_prefix['distributed'] + dist_job_prefix, dist_shards, dist_jobs = resolve_test_job_family( + periodic_wf, + "distributed", + "rocm", + dist_job_prefix, + rocm_shards["distributed"], + ) folder_list = get_or_create_test_folder(periodic_wf) # Download logs @@ -742,10 +786,6 @@ def main(): # HUD link: https://hud.pytorch.org/hud/pytorch/pytorch/main/1?per_page=50&name_filter=rocm # Make sure "Hide unstable jobs" is unselected, in case ROCm jobs are marked as unstable - if arch == "mi350": - dist_shards = 3 if not periodic_fallback_used else rocm_shards["distributed"] - else: - dist_shards = rocm_shards["distributed"] print(f"Using final ROCm shard count {dist_shards} for distributed") if not args.artifacts_only: @@ -753,7 +793,12 @@ def main(): [f"{current_prefix}rocm_dist{i}.txt", f"{dist_job_prefix} / test (distributed, {i}, {dist_shards}"] for i in range(1, dist_shards + 1) ] - download_logs(periodic_wf, test_log_list_rocm_distributed, folder_list[0]) + download_logs( + periodic_wf, + test_log_list_rocm_distributed, + folder_list[0], + jobs=dist_jobs, + ) # Download artifacts test_artifacts_list_rocm_distributed = [ @@ -800,23 +845,31 @@ def main(): default_wf_name = ROCmWorkflowNames['default'] if not default_fallback_used else default_fallbacks[arch][0] print(f"Using workflow '{default_wf_name}' with id:{rocm_wf['id']} for ROCm default{' (fallback)' if default_fallback_used else ''}") + default_prefix, default_shards, default_jobs = resolve_test_job_family( + rocm_wf, + "default", + "rocm", + rocm_job_prefix["default"], + rocm_shards["default"], + ) folder_list = get_or_create_test_folder(rocm_wf) # Download logs # If logs aren't found you might want to check the HUD for the correct tags # HUD link: https://hud.pytorch.org/hud/pytorch/pytorch/main/1?per_page=50&name_filter=rocm - if arch == "mi350": - default_shards = 6 if default_fallback_used else rocm_shards["default"] - else: - default_shards = rocm_shards["default"] print(f"Using final ROCm shard count {default_shards} for default") if not args.artifacts_only: test_log_list_rocm_default = [ - [f"{current_prefix}rocm{i}.txt", f"{rocm_job_prefix['default']} / test (default, {i}, {default_shards}"] + [f"{current_prefix}rocm{i}.txt", f"{default_prefix} / test (default, {i}, {default_shards}"] for i in range(1, default_shards + 1) ] - download_logs(rocm_wf, test_log_list_rocm_default, folder_list[0]) + download_logs( + rocm_wf, + test_log_list_rocm_default, + folder_list[0], + jobs=default_jobs, + ) # Download artifacts test_artifacts_list_rocm_default = [ @@ -842,10 +895,13 @@ def main(): print("===========================================") error_msg="Error: inductor workflow not found in scanned workflow runs. Try increasing max_pages." inductor_fallback_used = False - try: - inductor_wf_rocm = download_workflow_run(created=args.created, max_pages=args.max_pages, workflow=ROCmWorkflowNames["inductor"], sha=inductor_rocm_sha, ignore_status=args.ignore_status, status=status, error_msg=error_msg) - except (IndexError, Exception): - inductor_wf_rocm = None + if ROCmWorkflowNames["inductor"] == "trunk" and trunk_full_wf is not None: + inductor_wf_rocm = trunk_full_wf + else: + try: + inductor_wf_rocm = download_workflow_run(created=args.created, max_pages=args.max_pages, workflow=ROCmWorkflowNames["inductor"], sha=inductor_rocm_sha, ignore_status=args.ignore_status, status=status, error_msg=error_msg) + except (IndexError, Exception): + inductor_wf_rocm = None inductor_fallbacks = _fallbacks_for("inductor") if inductor_wf_rocm is None and arch in inductor_fallbacks: fallback_wf, fallback_prefix = inductor_fallbacks[arch] @@ -863,12 +919,19 @@ def main(): folder_list = get_or_create_test_folder(inductor_wf_rocm) - inductor_shards = rocm_shards["inductor"] - print(f"Using final ROCm shard count {inductor_shards} for inductor") if inductor_fallback_used and arch in inductor_fallbacks: inductor_job_prefix = inductor_fallbacks[arch][1] else: inductor_job_prefix = rocm_job_prefix['inductor'] + inductor_job_prefix, inductor_shards, inductor_jobs = \ + resolve_test_job_family( + inductor_wf_rocm, + "inductor", + "rocm", + inductor_job_prefix, + rocm_shards["inductor"], + ) + print(f"Using final ROCm shard count {inductor_shards} for inductor") # Download logs if not args.artifacts_only: @@ -876,7 +939,12 @@ def main(): [f"{current_prefix}rocm_inductor{i}.txt", f"{inductor_job_prefix} / test (inductor, {i}, {inductor_shards}"] for i in range(1, inductor_shards + 1) ] - download_logs(inductor_wf_rocm, test_log_list_rocm_inductor, folder_list[0]) + download_logs( + inductor_wf_rocm, + test_log_list_rocm_inductor, + folder_list[0], + jobs=inductor_jobs, + ) #Download artifacts test_artifacts_list_rocm_inductor = [ @@ -893,7 +961,7 @@ def main(): if not args.no_cuda: cuda_config = PARITY_CONFIG["cuda"] - cuda_job_prefix = CUDA_JOB_PREFIXES["default"] + cuda_job_prefix = trunk_cuda_job_prefix cuda_shards = cuda_config["shard_counts"] print("==========================================") print(f"Finding CUDA tests in workflow '{CUDAWorkflowNames['default']}' by sha: {sha}") @@ -925,8 +993,24 @@ def main(): folder_list = get_or_create_test_folder(trunk_wf) - cuda_default_shards = cuda_shards["default"] - cuda_dist_shards = cuda_shards["distributed"] + cuda_default_family = choose_test_job_family( + all_cuda_jobs, "default", "cuda", cuda_job_prefix + ) + cuda_dist_family = choose_test_job_family( + all_cuda_jobs, "distributed", "cuda", cuda_job_prefix + ) + cuda_default_shards = ( + cuda_default_family["total"] + if cuda_default_family else cuda_shards["default"] + ) + cuda_dist_shards = ( + cuda_dist_family["total"] + if cuda_dist_family else cuda_shards["distributed"] + ) + print( + f"Using final CUDA shard counts: default={cuda_default_shards}, " + f"distributed={cuda_dist_shards}" + ) # Download logs if not args.artifacts_only: @@ -1034,23 +1118,47 @@ def main(): baseline_xml_dir = os.path.join(test_folder, "baseline_xml") os.makedirs(baseline_xml_dir, exist_ok=True) + baseline_trunk_wf = None + if "trunk" in ROCmWorkflowNames.values(): + baseline_trunk_wf, _, _, _, _ = resolve_full_trunk_run( + baseline_sha, + CUDA_JOB_PREFIXES["default"], + status, + args.ignore_status, + ) + if not args.exclude_default: try: - baseline_default_wf = download_workflow_run( - created=args.created, max_pages=args.max_pages, - workflow=ROCmWorkflowNames["default"], sha=baseline_sha, - ignore_status=args.ignore_status, status=status, - error_msg=f"Baseline default workflow not found for {baseline_sha}", - ) + if ROCmWorkflowNames["default"] == "trunk" and baseline_trunk_wf: + baseline_default_wf = baseline_trunk_wf + else: + baseline_default_wf = download_workflow_run( + created=args.created, max_pages=args.max_pages, + workflow=ROCmWorkflowNames["default"], sha=baseline_sha, + ignore_status=args.ignore_status, status=status, + error_msg=f"Baseline default workflow not found for {baseline_sha}", + ) print(f"Baseline default workflow '{ROCmWorkflowNames['default']}' id: {baseline_default_wf['id']}") - default_shards = rocm_shards["default"] + baseline_default_prefix, default_shards, baseline_default_jobs = \ + resolve_test_job_family( + baseline_default_wf, + "default", + "rocm", + rocm_job_prefix["default"], + rocm_shards["default"], + ) if not args.artifacts_only: baseline_default_logs = [ - [f"{baseline_prefix}rocm{i}.txt", f"{rocm_job_prefix['default']} / test (default, {i}, {default_shards}"] + [f"{baseline_prefix}rocm{i}.txt", f"{baseline_default_prefix} / test (default, {i}, {default_shards}"] for i in range(1, default_shards + 1) ] - download_logs(baseline_default_wf, baseline_default_logs, test_folder) + download_logs( + baseline_default_wf, + baseline_default_logs, + test_folder, + jobs=baseline_default_jobs, + ) baseline_default_prefixes = [ f"test-reports-test-default-{i}-{default_shards}" @@ -1068,43 +1176,36 @@ def main(): if not args.exclude_distributed and "distributed" in ROCmWorkflowNames: try: - baseline_dist_workflow = ROCmWorkflowNames["distributed"] - baseline_dist_job_prefix = rocm_job_prefix["distributed"] - try: + if ROCmWorkflowNames["distributed"] == "trunk" and baseline_trunk_wf: + baseline_dist_wf = baseline_trunk_wf + else: baseline_dist_wf = download_workflow_run( created=args.created, max_pages=args.max_pages, - workflow=baseline_dist_workflow, sha=baseline_sha, + workflow=ROCmWorkflowNames["distributed"], sha=baseline_sha, ignore_status=args.ignore_status, status=status, error_msg=f"Baseline distributed workflow not found for {baseline_sha}", ) - except (IndexError, Exception): - baseline_dist_wf = None - - baseline_dist_fallbacks = _fallbacks_for("distributed") - if baseline_dist_wf is None and arch in baseline_dist_fallbacks: - baseline_dist_workflow, baseline_dist_job_prefix = \ - baseline_dist_fallbacks[arch] - print( - f"Baseline distributed not found in " - f"{ROCmWorkflowNames['distributed']}, falling back to " - f"{baseline_dist_workflow}" + print(f"Baseline distributed workflow '{ROCmWorkflowNames['distributed']}' id: {baseline_dist_wf['id']}") + baseline_dist_prefix, dist_shards, baseline_dist_jobs = \ + resolve_test_job_family( + baseline_dist_wf, + "distributed", + "rocm", + rocm_job_prefix["distributed"], + rocm_shards["distributed"], ) - baseline_dist_wf = download_workflow_run( - created=args.created, max_pages=args.max_pages, - workflow=baseline_dist_workflow, sha=baseline_sha, - ignore_status=args.ignore_status, status=status, - error_msg=f"Baseline distributed fallback not found for {baseline_sha}", - ) - - print(f"Baseline distributed workflow '{baseline_dist_workflow}' id: {baseline_dist_wf['id']}") - dist_shards = rocm_shards["distributed"] if not args.artifacts_only: baseline_dist_logs = [ - [f"{baseline_prefix}rocm_dist{i}.txt", f"{baseline_dist_job_prefix} / test (distributed, {i}, {dist_shards}"] + [f"{baseline_prefix}rocm_dist{i}.txt", f"{baseline_dist_prefix} / test (distributed, {i}, {dist_shards}"] for i in range(1, dist_shards + 1) ] - download_logs(baseline_dist_wf, baseline_dist_logs, test_folder) + download_logs( + baseline_dist_wf, + baseline_dist_logs, + test_folder, + jobs=baseline_dist_jobs, + ) baseline_dist_prefixes = [ f"test-reports-test-distributed-{i}-{dist_shards}" @@ -1122,21 +1223,36 @@ def main(): if not args.exclude_inductor and "inductor" in ROCmWorkflowNames: try: - baseline_inductor_wf = download_workflow_run( - created=args.created, max_pages=args.max_pages, - workflow=ROCmWorkflowNames["inductor"], sha=baseline_sha, - ignore_status=args.ignore_status, status=status, - error_msg=f"Baseline inductor workflow not found for {baseline_sha}", - ) + if ROCmWorkflowNames["inductor"] == "trunk" and baseline_trunk_wf: + baseline_inductor_wf = baseline_trunk_wf + else: + baseline_inductor_wf = download_workflow_run( + created=args.created, max_pages=args.max_pages, + workflow=ROCmWorkflowNames["inductor"], sha=baseline_sha, + ignore_status=args.ignore_status, status=status, + error_msg=f"Baseline inductor workflow not found for {baseline_sha}", + ) print(f"Baseline inductor workflow '{ROCmWorkflowNames['inductor']}' id: {baseline_inductor_wf['id']}") - inductor_shards = rocm_shards["inductor"] + baseline_inductor_prefix, inductor_shards, baseline_inductor_jobs = \ + resolve_test_job_family( + baseline_inductor_wf, + "inductor", + "rocm", + rocm_job_prefix["inductor"], + rocm_shards["inductor"], + ) if not args.artifacts_only: baseline_inductor_logs = [ - [f"{baseline_prefix}rocm_inductor{i}.txt", f"{rocm_job_prefix['inductor']} / test (inductor, {i}, {inductor_shards}"] + [f"{baseline_prefix}rocm_inductor{i}.txt", f"{baseline_inductor_prefix} / test (inductor, {i}, {inductor_shards}"] for i in range(1, inductor_shards + 1) ] - download_logs(baseline_inductor_wf, baseline_inductor_logs, test_folder) + download_logs( + baseline_inductor_wf, + baseline_inductor_logs, + test_folder, + jobs=baseline_inductor_jobs, + ) baseline_inductor_prefixes = [ f"test-reports-test-inductor-{i}-{inductor_shards}" diff --git a/.automation_scripts/pytorch-unit-test-scripts/job_name_match.py b/.automation_scripts/pytorch-unit-test-scripts/job_name_match.py new file mode 100644 index 0000000000000..aa60bbbc144ca --- /dev/null +++ b/.automation_scripts/pytorch-unit-test-scripts/job_name_match.py @@ -0,0 +1,58 @@ +import difflib +import re + + +_TEST_JOB = re.compile( + r"^(?P.+) / (?Ptest(?:-osdc)?) " + r"\((?P[^,]+), (?P\d+), (?P\d+)" +) +_SPECIAL_FAMILIES = ("build-only", "debug", "no-ops", "slow-gradcheck", "smoke") + + +def _matches_platform(prefix, platform): + normalized = prefix.lower() + if platform == "rocm": + return "rocm" in normalized and "cuda" not in normalized + if platform == "cuda": + return "cuda" in normalized and "rocm" not in normalized + raise ValueError(f"Unsupported platform: {platform}") + + +def choose_test_job_family(jobs, test_config, platform, configured_prefix): + """Return the best matching sharded test family from actual workflow jobs.""" + families = {} + for job in jobs: + match = _TEST_JOB.match(job.get("name", "")) + if not match or match.group("config") != test_config: + continue + prefix = match.group("prefix") + if not _matches_platform(prefix, platform): + continue + key = (prefix, match.group("kind"), int(match.group("total"))) + families.setdefault(key, []).append(job) + + if not families: + return None + + def score(item): + (prefix, _, total), matching_jobs = item + shards = { + int(_TEST_JOB.match(job["name"]).group("shard")) + for job in matching_jobs + } + normalized = prefix.lower() + return ( + prefix == configured_prefix, + not any(marker in normalized for marker in _SPECIAL_FAMILIES), + len(shards) == total, + difflib.SequenceMatcher(None, prefix, configured_prefix).ratio(), + len(shards), + ) + + (prefix, kind, total), matching_jobs = max(families.items(), key=score) + return { + "prefix": prefix, + "kind": kind, + "total": total, + "jobs": matching_jobs, + } diff --git a/.automation_scripts/pytorch-unit-test-scripts/parity_job_config.json b/.automation_scripts/pytorch-unit-test-scripts/parity_job_config.json index 3e5d6823137b7..6c6c38254e6ca 100644 --- a/.automation_scripts/pytorch-unit-test-scripts/parity_job_config.json +++ b/.automation_scripts/pytorch-unit-test-scripts/parity_job_config.json @@ -20,17 +20,9 @@ }, "rocm": { "mi350": { - "default": [ - { "workflow": "trunk", "job_prefix": "linux-noble-rocm-py3.11-mi350" }, - { "workflow": "rocm-mi350", "job_prefix": "linux-noble-rocm-py3.12-mi350" } - ], - "distributed": [ - { "workflow": "periodic-rocm-mi350", "job_prefix": "linux-noble-rocm-py3.12-mi350" }, - { "workflow": "trunk", "job_prefix": "linux-noble-rocm-py3.11-mi350" } - ], - "inductor": [ - { "workflow": "trunk", "job_prefix": "linux-noble-rocm-py3.11-mi350" } - ], + "default": [{ "workflow": "trunk", "job_prefix": "linux-noble-rocm-py3.11-mi350" }], + "distributed": [{ "workflow": "trunk", "job_prefix": "linux-noble-rocm-py3.11-mi350" }], + "inductor": [{ "workflow": "trunk", "job_prefix": "linux-noble-rocm-py3.11-mi350" }], "shard_counts": { "default": 8, "distributed": 3, "inductor": 2 }, "checkrun_regex": "rocm.*mi350.*/ test [(](default|distributed|inductor)," }, diff --git a/.automation_scripts/pytorch-unit-test-scripts/test_job_name_match.py b/.automation_scripts/pytorch-unit-test-scripts/test_job_name_match.py new file mode 100644 index 0000000000000..23fa756eee5e5 --- /dev/null +++ b/.automation_scripts/pytorch-unit-test-scripts/test_job_name_match.py @@ -0,0 +1,98 @@ +import unittest + +from job_name_match import choose_test_job_family + + +def job(name, job_id): + return {"name": name, "id": job_id} + + +class ChooseTestJobFamilyTest(unittest.TestCase): + def test_discovers_renamed_cuda_family(self): + jobs = [ + job( + f"linux-jammy-cuda13.2-py3.10-gcc11 / test " + f"(default, {shard}, 14, runner)", + shard, + ) + for shard in range(1, 15) + ] + + family = choose_test_job_family( + jobs, + "default", + "cuda", + "linux-jammy-cuda13.0-py3.10-gcc11", + ) + + self.assertEqual(family["prefix"], "linux-jammy-cuda13.2-py3.10-gcc11") + self.assertEqual(family["total"], 14) + + def test_prefers_normal_cuda_family_over_debug(self): + jobs = [ + job( + f"linux-jammy-cuda13.2-py3.10-gcc11 / test " + f"(default, {shard}, 14, runner)", + shard, + ) + for shard in range(1, 15) + ] + jobs.extend( + job( + f"linux-jammy-cuda13.0-py3.10-gcc11-debug / test " + f"(default, {shard}, 7, runner)", + 100 + shard, + ) + for shard in range(1, 8) + ) + + family = choose_test_job_family( + jobs, + "default", + "cuda", + "linux-jammy-cuda13.0-py3.10-gcc11", + ) + + self.assertEqual(family["prefix"], "linux-jammy-cuda13.2-py3.10-gcc11") + + def test_discovers_historical_rocm_family(self): + jobs = [ + job( + f"linux-jammy-rocm-py3.10-mi350 / test " + f"(distributed, {shard}, 3, runner)", + shard, + ) + for shard in range(1, 4) + ] + + family = choose_test_job_family( + jobs, + "distributed", + "rocm", + "linux-noble-rocm-py3.11-mi350", + ) + + self.assertEqual(family["prefix"], "linux-jammy-rocm-py3.10-mi350") + self.assertEqual(family["total"], 3) + + def test_keeps_platforms_separate(self): + jobs = [ + job( + "linux-jammy-cuda13.2-py3.10-gcc11 / test " + "(default, 1, 14, runner)", + 1, + ) + ] + + self.assertIsNone( + choose_test_job_family( + jobs, + "default", + "rocm", + "linux-noble-rocm-py3.11-mi350", + ) + ) + + +if __name__ == "__main__": + unittest.main() From d5a5cd46ed41365c9352120aea9a1b7c60f2eaa8 Mon Sep 17 00:00:00 2001 From: ethanwee1 Date: Wed, 16 Sep 2026 19:32:42 +0000 Subject: [PATCH 2/3] [CI] Harden parity job family selection Scope fallback checks to trunk push runs, require complete shard families, preserve discovered test kinds, and retain configured baseline fallbacks for other architectures. --- .../download_testlogs | 239 +++++++++------- .../job_name_match.py | 122 +++++---- .../test_job_name_match.py | 259 +++++++++++------- 3 files changed, 364 insertions(+), 256 deletions(-) diff --git a/.automation_scripts/pytorch-unit-test-scripts/download_testlogs b/.automation_scripts/pytorch-unit-test-scripts/download_testlogs index c21da09fdd33a..d77dd2d5e8185 100755 --- a/.automation_scripts/pytorch-unit-test-scripts/download_testlogs +++ b/.automation_scripts/pytorch-unit-test-scripts/download_testlogs @@ -242,13 +242,13 @@ def resolve_test_job_family(wf, test_config, platform, configured_prefix, jobs, test_config, platform, configured_prefix ) if family is None: - return configured_prefix, fallback_shards, jobs + return configured_prefix, fallback_shards, "test", jobs if family["prefix"] != configured_prefix: print( f"NOTE: configured {platform} prefix '{configured_prefix}' is absent; " f"using '{family['prefix']}' from run {wf['id']}" ) - return family["prefix"], family["total"], jobs + return family["prefix"], family["total"], family["kind"], jobs def get_cuda_inductor_test_jobs(jobs): """Return the CUDA inductor test kind and jobs for either main or PR CI layouts.""" @@ -574,59 +574,54 @@ def resolve_full_trunk_run(sha, cuda_job_prefix, status='success', ignore_status headers=authentication_headers, params=params, ) trunk_runs = resp.json().get('workflow_runs', []) - # Prefer the normal push trunk run over scheduled trunk runs. The scheduled - # trunk.yml runs are the periodic rerun_disabled_tests / mem_leak_check - # variants: their test jobs share the same names but run in a degenerate - # mode and produce near-empty reports. A stable sort keeps the API's - # newest-first order within each group, so we pick the newest push run. - trunk_runs.sort(key=lambda r: r.get('event') != 'push') + # Scheduled trunk.yml runs are periodic rerun_disabled_tests / mem_leak_check + # variants. Never let those preempt a normal push for the same SHA. + push_runs = [run for run in trunk_runs if run.get("event") == "push"] + candidate_runs = push_runs or trunk_runs trunk_wf = None cuda_test_jobs = [] cuda_test_kind = CUDA_TEST_KINDS[0] all_cuda_jobs = [] - for run in trunk_runs: + for run in candidate_runs: jobs = get_workflow_jobs(run) family = choose_test_job_family( jobs, "default", "cuda", cuda_job_prefix ) - if family: + if family and family["complete"]: trunk_wf = run all_cuda_jobs = jobs - cuda_test_jobs = family["jobs"] - cuda_test_kind = family["kind"] cuda_job_prefix = family["prefix"] + cuda_test_kind, cuda_test_jobs = get_cuda_test_jobs( + jobs, cuda_job_prefix + ) print(f"Found CUDA test jobs in trunk run {run['id']}") break - if not cuda_test_jobs and trunk_runs: - # CUDA jobs may live in a run the jobs API does not surface (e.g. a - # reusable workflow); the check-runs API resolves the actual run. - print("No CUDA test jobs in any trunk run's jobs API, trying check-runs API...") + if not cuda_test_jobs and candidate_runs: + # Check runs can include periodic and other workflows for the same SHA. + # Scope the fallback to the candidate trunk runs before choosing a family. + print("No CUDA test jobs in trunk jobs APIs, trying scoped check runs...") check_runs = get_check_runs_for_commit(sha, "") - family = choose_test_job_family( - check_runs, "default", "cuda", cuda_job_prefix - ) - if family: - cuda_test_kind = family["kind"] - cuda_test_jobs = family["jobs"] - cuda_job_prefix = family["prefix"] - run_match = re.search(r'/runs/(\d+)/', cuda_test_jobs[0].get('details_url', '')) - if run_match: - actual_run_id = int(run_match.group(1)) - trunk_wf = next((r for r in trunk_runs if r['id'] == actual_run_id), None) - if trunk_wf is None: - resp = requests.get( - f"https://api.github.com/repos/pytorch/pytorch/actions/runs/{actual_run_id}", - headers=authentication_headers, - ) - trunk_wf = resp.json() + for run in candidate_runs: + marker = f"/runs/{run['id']}/" + run_checks = [ + check for check in check_runs + if marker in (check.get("details_url") or "") + ] + family = choose_test_job_family( + run_checks, "default", "cuda", cuda_job_prefix + ) + if family and family["complete"]: + trunk_wf = run + all_cuda_jobs = run_checks + cuda_job_prefix = family["prefix"] + cuda_test_kind, cuda_test_jobs = get_cuda_test_jobs( + run_checks, cuda_job_prefix + ) print(f"CUDA test jobs are in trunk run {trunk_wf['id']} (found via check-runs)") - all_cuda_jobs = list(cuda_test_jobs) - - if trunk_wf is None and trunk_runs: - trunk_wf = trunk_runs[0] + break return ( trunk_wf, @@ -772,7 +767,8 @@ def main(): else: dist_job_prefix = rocm_job_prefix['distributed'] - dist_job_prefix, dist_shards, dist_jobs = resolve_test_job_family( + dist_job_prefix, dist_shards, dist_job_kind, dist_jobs = \ + resolve_test_job_family( periodic_wf, "distributed", "rocm", @@ -790,7 +786,7 @@ def main(): if not args.artifacts_only: test_log_list_rocm_distributed = [ - [f"{current_prefix}rocm_dist{i}.txt", f"{dist_job_prefix} / test (distributed, {i}, {dist_shards}"] + [f"{current_prefix}rocm_dist{i}.txt", f"{dist_job_prefix} / {dist_job_kind} (distributed, {i}, {dist_shards}"] for i in range(1, dist_shards + 1) ] download_logs( @@ -802,7 +798,7 @@ def main(): # Download artifacts test_artifacts_list_rocm_distributed = [ - f"test-reports-test-distributed-{i}-{dist_shards}" + f"test-reports-{dist_job_kind}-distributed-{i}-{dist_shards}" for i in range(1, dist_shards + 1) ] download_artifacts( @@ -845,13 +841,14 @@ def main(): default_wf_name = ROCmWorkflowNames['default'] if not default_fallback_used else default_fallbacks[arch][0] print(f"Using workflow '{default_wf_name}' with id:{rocm_wf['id']} for ROCm default{' (fallback)' if default_fallback_used else ''}") - default_prefix, default_shards, default_jobs = resolve_test_job_family( - rocm_wf, - "default", - "rocm", - rocm_job_prefix["default"], - rocm_shards["default"], - ) + default_prefix, default_shards, default_job_kind, default_jobs = \ + resolve_test_job_family( + rocm_wf, + "default", + "rocm", + rocm_job_prefix["default"], + rocm_shards["default"], + ) folder_list = get_or_create_test_folder(rocm_wf) # Download logs @@ -861,7 +858,7 @@ def main(): if not args.artifacts_only: test_log_list_rocm_default = [ - [f"{current_prefix}rocm{i}.txt", f"{default_prefix} / test (default, {i}, {default_shards}"] + [f"{current_prefix}rocm{i}.txt", f"{default_prefix} / {default_job_kind} (default, {i}, {default_shards}"] for i in range(1, default_shards + 1) ] download_logs( @@ -873,7 +870,7 @@ def main(): # Download artifacts test_artifacts_list_rocm_default = [ - f"test-reports-test-default-{i}-{default_shards}" + f"test-reports-{default_job_kind}-default-{i}-{default_shards}" for i in range(1, default_shards + 1) ] if not args.exclude_default: @@ -923,7 +920,8 @@ def main(): inductor_job_prefix = inductor_fallbacks[arch][1] else: inductor_job_prefix = rocm_job_prefix['inductor'] - inductor_job_prefix, inductor_shards, inductor_jobs = \ + inductor_job_prefix, inductor_shards, inductor_job_kind, \ + inductor_jobs = \ resolve_test_job_family( inductor_wf_rocm, "inductor", @@ -936,7 +934,7 @@ def main(): # Download logs if not args.artifacts_only: test_log_list_rocm_inductor = [ - [f"{current_prefix}rocm_inductor{i}.txt", f"{inductor_job_prefix} / test (inductor, {i}, {inductor_shards}"] + [f"{current_prefix}rocm_inductor{i}.txt", f"{inductor_job_prefix} / {inductor_job_kind} (inductor, {i}, {inductor_shards}"] for i in range(1, inductor_shards + 1) ] download_logs( @@ -948,7 +946,7 @@ def main(): #Download artifacts test_artifacts_list_rocm_inductor = [ - f"test-reports-test-inductor-{i}-{inductor_shards}" + f"test-reports-{inductor_job_kind}-inductor-{i}-{inductor_shards}" for i in range(1, inductor_shards + 1) ] download_artifacts( @@ -1003,10 +1001,18 @@ def main(): cuda_default_family["total"] if cuda_default_family else cuda_shards["default"] ) + cuda_default_kind = ( + cuda_default_family["kind"] + if cuda_default_family else cuda_test_job_kind + ) cuda_dist_shards = ( cuda_dist_family["total"] if cuda_dist_family else cuda_shards["distributed"] ) + cuda_dist_kind = ( + cuda_dist_family["kind"] + if cuda_dist_family else cuda_test_job_kind + ) print( f"Using final CUDA shard counts: default={cuda_default_shards}, " f"distributed={cuda_dist_shards}" @@ -1015,12 +1021,12 @@ def main(): # Download logs if not args.artifacts_only: test_log_list_cuda = [ - [f"cuda{i}.txt", f"{cuda_job_prefix} / {cuda_test_job_kind} (default, {i}, {cuda_default_shards}"] + [f"cuda{i}.txt", f"{cuda_job_prefix} / {cuda_default_kind} (default, {i}, {cuda_default_shards}"] for i in range(1, cuda_default_shards + 1) ] if not args.exclude_distributed: test_log_list_cuda += [ - [f"cuda_dist{i}.txt", f"{cuda_job_prefix} / {cuda_test_job_kind} (distributed, {i}, {cuda_dist_shards}"] + [f"cuda_dist{i}.txt", f"{cuda_job_prefix} / {cuda_dist_kind} (distributed, {i}, {cuda_dist_shards}"] for i in range(1, cuda_dist_shards + 1) ] @@ -1030,13 +1036,13 @@ def main(): test_artifacts_list_cuda = [] if not args.exclude_default: test_artifacts_list_cuda += [ - f"test-reports-{cuda_test_job_kind}-default-{i}-{cuda_default_shards}" + f"test-reports-{cuda_default_kind}-default-{i}-{cuda_default_shards}" for i in range(1, cuda_default_shards + 1) ] if not args.exclude_distributed: test_artifacts_list_cuda += [ - f"test-reports-{cuda_test_job_kind}-distributed-{i}-{cuda_dist_shards}" + f"test-reports-{cuda_dist_kind}-distributed-{i}-{cuda_dist_shards}" for i in range(1, cuda_dist_shards + 1) ] @@ -1127,30 +1133,74 @@ def main(): args.ignore_status, ) - if not args.exclude_default: + def resolve_baseline_source(test_config): + workflow = ROCmWorkflowNames[test_config] + job_prefix = rocm_job_prefix[test_config] + if workflow == "trunk" and baseline_trunk_wf: + return baseline_trunk_wf, workflow, job_prefix try: - if ROCmWorkflowNames["default"] == "trunk" and baseline_trunk_wf: - baseline_default_wf = baseline_trunk_wf + workflow_run = download_workflow_run( + created=args.created, + max_pages=args.max_pages, + workflow=workflow, + sha=baseline_sha, + ignore_status=args.ignore_status, + status=status, + error_msg=( + f"Baseline {test_config} workflow not found for " + f"{baseline_sha}" + ), + ) + except (IndexError, Exception): + workflow_run = None + + fallbacks = _fallbacks_for(test_config) + if workflow_run is None and arch in fallbacks: + workflow, job_prefix = fallbacks[arch] + print( + f"Baseline {test_config} not found in " + f"{ROCmWorkflowNames[test_config]}, falling back to {workflow}" + ) + if workflow == "trunk" and baseline_trunk_wf: + workflow_run = baseline_trunk_wf else: - baseline_default_wf = download_workflow_run( - created=args.created, max_pages=args.max_pages, - workflow=ROCmWorkflowNames["default"], sha=baseline_sha, - ignore_status=args.ignore_status, status=status, - error_msg=f"Baseline default workflow not found for {baseline_sha}", + workflow_run = download_workflow_run( + created=args.created, + max_pages=args.max_pages, + workflow=workflow, + sha=baseline_sha, + ignore_status=args.ignore_status, + status=status, + error_msg=( + f"Baseline {test_config} fallback not found for " + f"{baseline_sha}" + ), ) - print(f"Baseline default workflow '{ROCmWorkflowNames['default']}' id: {baseline_default_wf['id']}") - baseline_default_prefix, default_shards, baseline_default_jobs = \ - resolve_test_job_family( + if workflow_run is None: + raise Exception( + f"Baseline {test_config} workflow not found for {baseline_sha}" + ) + return workflow_run, workflow, job_prefix + + if not args.exclude_default: + try: + baseline_default_wf, baseline_default_workflow, \ + baseline_default_configured_prefix = \ + resolve_baseline_source("default") + print(f"Baseline default workflow '{baseline_default_workflow}' id: {baseline_default_wf['id']}") + baseline_default_prefix, default_shards, \ + baseline_default_kind, baseline_default_jobs = \ + resolve_test_job_family( baseline_default_wf, "default", "rocm", - rocm_job_prefix["default"], + baseline_default_configured_prefix, rocm_shards["default"], ) if not args.artifacts_only: baseline_default_logs = [ - [f"{baseline_prefix}rocm{i}.txt", f"{baseline_default_prefix} / test (default, {i}, {default_shards}"] + [f"{baseline_prefix}rocm{i}.txt", f"{baseline_default_prefix} / {baseline_default_kind} (default, {i}, {default_shards}"] for i in range(1, default_shards + 1) ] download_logs( @@ -1161,7 +1211,7 @@ def main(): ) baseline_default_prefixes = [ - f"test-reports-test-default-{i}-{default_shards}" + f"test-reports-{baseline_default_kind}-default-{i}-{default_shards}" for i in range(1, default_shards + 1) ] download_artifacts( @@ -1176,28 +1226,22 @@ def main(): if not args.exclude_distributed and "distributed" in ROCmWorkflowNames: try: - if ROCmWorkflowNames["distributed"] == "trunk" and baseline_trunk_wf: - baseline_dist_wf = baseline_trunk_wf - else: - baseline_dist_wf = download_workflow_run( - created=args.created, max_pages=args.max_pages, - workflow=ROCmWorkflowNames["distributed"], sha=baseline_sha, - ignore_status=args.ignore_status, status=status, - error_msg=f"Baseline distributed workflow not found for {baseline_sha}", - ) - print(f"Baseline distributed workflow '{ROCmWorkflowNames['distributed']}' id: {baseline_dist_wf['id']}") - baseline_dist_prefix, dist_shards, baseline_dist_jobs = \ - resolve_test_job_family( + baseline_dist_wf, baseline_dist_workflow, \ + baseline_dist_configured_prefix = \ + resolve_baseline_source("distributed") + print(f"Baseline distributed workflow '{baseline_dist_workflow}' id: {baseline_dist_wf['id']}") + baseline_dist_prefix, dist_shards, baseline_dist_kind, \ + baseline_dist_jobs = resolve_test_job_family( baseline_dist_wf, "distributed", "rocm", - rocm_job_prefix["distributed"], + baseline_dist_configured_prefix, rocm_shards["distributed"], ) if not args.artifacts_only: baseline_dist_logs = [ - [f"{baseline_prefix}rocm_dist{i}.txt", f"{baseline_dist_prefix} / test (distributed, {i}, {dist_shards}"] + [f"{baseline_prefix}rocm_dist{i}.txt", f"{baseline_dist_prefix} / {baseline_dist_kind} (distributed, {i}, {dist_shards}"] for i in range(1, dist_shards + 1) ] download_logs( @@ -1208,7 +1252,7 @@ def main(): ) baseline_dist_prefixes = [ - f"test-reports-test-distributed-{i}-{dist_shards}" + f"test-reports-{baseline_dist_kind}-distributed-{i}-{dist_shards}" for i in range(1, dist_shards + 1) ] download_artifacts( @@ -1223,28 +1267,23 @@ def main(): if not args.exclude_inductor and "inductor" in ROCmWorkflowNames: try: - if ROCmWorkflowNames["inductor"] == "trunk" and baseline_trunk_wf: - baseline_inductor_wf = baseline_trunk_wf - else: - baseline_inductor_wf = download_workflow_run( - created=args.created, max_pages=args.max_pages, - workflow=ROCmWorkflowNames["inductor"], sha=baseline_sha, - ignore_status=args.ignore_status, status=status, - error_msg=f"Baseline inductor workflow not found for {baseline_sha}", - ) - print(f"Baseline inductor workflow '{ROCmWorkflowNames['inductor']}' id: {baseline_inductor_wf['id']}") - baseline_inductor_prefix, inductor_shards, baseline_inductor_jobs = \ - resolve_test_job_family( + baseline_inductor_wf, baseline_inductor_workflow, \ + baseline_inductor_configured_prefix = \ + resolve_baseline_source("inductor") + print(f"Baseline inductor workflow '{baseline_inductor_workflow}' id: {baseline_inductor_wf['id']}") + baseline_inductor_prefix, inductor_shards, \ + baseline_inductor_kind, baseline_inductor_jobs = \ + resolve_test_job_family( baseline_inductor_wf, "inductor", "rocm", - rocm_job_prefix["inductor"], + baseline_inductor_configured_prefix, rocm_shards["inductor"], ) if not args.artifacts_only: baseline_inductor_logs = [ - [f"{baseline_prefix}rocm_inductor{i}.txt", f"{baseline_inductor_prefix} / test (inductor, {i}, {inductor_shards}"] + [f"{baseline_prefix}rocm_inductor{i}.txt", f"{baseline_inductor_prefix} / {baseline_inductor_kind} (inductor, {i}, {inductor_shards}"] for i in range(1, inductor_shards + 1) ] download_logs( @@ -1255,7 +1294,7 @@ def main(): ) baseline_inductor_prefixes = [ - f"test-reports-test-inductor-{i}-{inductor_shards}" + f"test-reports-{baseline_inductor_kind}-inductor-{i}-{inductor_shards}" for i in range(1, inductor_shards + 1) ] download_artifacts( diff --git a/.automation_scripts/pytorch-unit-test-scripts/job_name_match.py b/.automation_scripts/pytorch-unit-test-scripts/job_name_match.py index aa60bbbc144ca..401ff0fa10dfa 100644 --- a/.automation_scripts/pytorch-unit-test-scripts/job_name_match.py +++ b/.automation_scripts/pytorch-unit-test-scripts/job_name_match.py @@ -1,58 +1,64 @@ -import difflib -import re - - -_TEST_JOB = re.compile( - r"^(?P.+) / (?Ptest(?:-osdc)?) " - r"\((?P[^,]+), (?P\d+), (?P\d+)" -) -_SPECIAL_FAMILIES = ("build-only", "debug", "no-ops", "slow-gradcheck", "smoke") - - -def _matches_platform(prefix, platform): - normalized = prefix.lower() - if platform == "rocm": - return "rocm" in normalized and "cuda" not in normalized - if platform == "cuda": - return "cuda" in normalized and "rocm" not in normalized - raise ValueError(f"Unsupported platform: {platform}") - - -def choose_test_job_family(jobs, test_config, platform, configured_prefix): - """Return the best matching sharded test family from actual workflow jobs.""" - families = {} - for job in jobs: - match = _TEST_JOB.match(job.get("name", "")) - if not match or match.group("config") != test_config: - continue - prefix = match.group("prefix") - if not _matches_platform(prefix, platform): - continue - key = (prefix, match.group("kind"), int(match.group("total"))) - families.setdefault(key, []).append(job) - - if not families: - return None - - def score(item): - (prefix, _, total), matching_jobs = item - shards = { - int(_TEST_JOB.match(job["name"]).group("shard")) - for job in matching_jobs - } - normalized = prefix.lower() - return ( - prefix == configured_prefix, - not any(marker in normalized for marker in _SPECIAL_FAMILIES), - len(shards) == total, - difflib.SequenceMatcher(None, prefix, configured_prefix).ratio(), - len(shards), - ) - - (prefix, kind, total), matching_jobs = max(families.items(), key=score) - return { - "prefix": prefix, - "kind": kind, - "total": total, - "jobs": matching_jobs, - } +import difflib +import re + + +_TEST_JOB = re.compile( + r"^(?P.+) / (?Ptest(?:-osdc)?) " + r"\((?P[^,]+), (?P\d+), (?P\d+)" +) +_SPECIAL_FAMILIES = ("build-only", "debug", "no-ops", "slow-gradcheck", "smoke") + + +def _matches_platform(prefix, platform): + normalized = prefix.lower() + if platform == "rocm": + return "rocm" in normalized and "cuda" not in normalized + if platform == "cuda": + return "cuda" in normalized and "rocm" not in normalized + raise ValueError(f"Unsupported platform: {platform}") + + +def choose_test_job_family(jobs, test_config, platform, configured_prefix): + """Return the best matching sharded test family from actual workflow jobs.""" + families = {} + for job in jobs: + match = _TEST_JOB.match(job.get("name", "")) + if not match or match.group("config") != test_config: + continue + prefix = match.group("prefix") + if not _matches_platform(prefix, platform): + continue + key = (prefix, match.group("kind"), int(match.group("total"))) + families.setdefault(key, []).append(job) + + if not families: + return None + + def score(item): + (prefix, _, total), matching_jobs = item + shards = { + int(_TEST_JOB.match(job["name"]).group("shard")) + for job in matching_jobs + } + complete = shards == set(range(1, total + 1)) + normalized = prefix.lower() + return ( + not any(marker in normalized for marker in _SPECIAL_FAMILIES), + complete, + prefix == configured_prefix, + difflib.SequenceMatcher(None, prefix, configured_prefix).ratio(), + len(shards), + ) + + (prefix, kind, total), matching_jobs = max(families.items(), key=score) + shards = { + int(_TEST_JOB.match(job["name"]).group("shard")) + for job in matching_jobs + } + return { + "prefix": prefix, + "kind": kind, + "total": total, + "jobs": matching_jobs, + "complete": shards == set(range(1, total + 1)), + } diff --git a/.automation_scripts/pytorch-unit-test-scripts/test_job_name_match.py b/.automation_scripts/pytorch-unit-test-scripts/test_job_name_match.py index 23fa756eee5e5..3df30aa8016b3 100644 --- a/.automation_scripts/pytorch-unit-test-scripts/test_job_name_match.py +++ b/.automation_scripts/pytorch-unit-test-scripts/test_job_name_match.py @@ -1,98 +1,161 @@ -import unittest - -from job_name_match import choose_test_job_family - - -def job(name, job_id): - return {"name": name, "id": job_id} - - -class ChooseTestJobFamilyTest(unittest.TestCase): - def test_discovers_renamed_cuda_family(self): - jobs = [ - job( - f"linux-jammy-cuda13.2-py3.10-gcc11 / test " - f"(default, {shard}, 14, runner)", - shard, - ) - for shard in range(1, 15) - ] - - family = choose_test_job_family( - jobs, - "default", - "cuda", - "linux-jammy-cuda13.0-py3.10-gcc11", - ) - - self.assertEqual(family["prefix"], "linux-jammy-cuda13.2-py3.10-gcc11") - self.assertEqual(family["total"], 14) - - def test_prefers_normal_cuda_family_over_debug(self): - jobs = [ - job( - f"linux-jammy-cuda13.2-py3.10-gcc11 / test " - f"(default, {shard}, 14, runner)", - shard, - ) - for shard in range(1, 15) - ] - jobs.extend( - job( - f"linux-jammy-cuda13.0-py3.10-gcc11-debug / test " - f"(default, {shard}, 7, runner)", - 100 + shard, - ) - for shard in range(1, 8) - ) - - family = choose_test_job_family( - jobs, - "default", - "cuda", - "linux-jammy-cuda13.0-py3.10-gcc11", - ) - - self.assertEqual(family["prefix"], "linux-jammy-cuda13.2-py3.10-gcc11") - - def test_discovers_historical_rocm_family(self): - jobs = [ - job( - f"linux-jammy-rocm-py3.10-mi350 / test " - f"(distributed, {shard}, 3, runner)", - shard, - ) - for shard in range(1, 4) - ] - - family = choose_test_job_family( - jobs, - "distributed", - "rocm", - "linux-noble-rocm-py3.11-mi350", - ) - - self.assertEqual(family["prefix"], "linux-jammy-rocm-py3.10-mi350") - self.assertEqual(family["total"], 3) - - def test_keeps_platforms_separate(self): - jobs = [ - job( - "linux-jammy-cuda13.2-py3.10-gcc11 / test " - "(default, 1, 14, runner)", - 1, - ) - ] - - self.assertIsNone( - choose_test_job_family( - jobs, - "default", - "rocm", - "linux-noble-rocm-py3.11-mi350", - ) - ) - - -if __name__ == "__main__": - unittest.main() +import unittest + +from job_name_match import choose_test_job_family + + +def job(name, job_id): + return {"name": name, "id": job_id} + + +class ChooseTestJobFamilyTest(unittest.TestCase): + def test_discovers_renamed_cuda_family(self): + jobs = [ + job( + f"linux-jammy-cuda13.2-py3.10-gcc11 / test " + f"(default, {shard}, 14, runner)", + shard, + ) + for shard in range(1, 15) + ] + + family = choose_test_job_family( + jobs, + "default", + "cuda", + "linux-jammy-cuda13.0-py3.10-gcc11", + ) + + self.assertEqual(family["prefix"], "linux-jammy-cuda13.2-py3.10-gcc11") + self.assertEqual(family["total"], 14) + + def test_prefers_normal_cuda_family_over_debug(self): + jobs = [ + job( + f"linux-jammy-cuda13.2-py3.10-gcc11 / test " + f"(default, {shard}, 14, runner)", + shard, + ) + for shard in range(1, 15) + ] + jobs.extend( + job( + f"linux-jammy-cuda13.0-py3.10-gcc11-debug / test " + f"(default, {shard}, 7, runner)", + 100 + shard, + ) + for shard in range(1, 8) + ) + + family = choose_test_job_family( + jobs, + "default", + "cuda", + "linux-jammy-cuda13.0-py3.10-gcc11", + ) + + self.assertEqual(family["prefix"], "linux-jammy-cuda13.2-py3.10-gcc11") + + def test_prefers_complete_renamed_family_over_incomplete_exact_match(self): + jobs = [ + job( + "linux-jammy-cuda13.0-py3.10-gcc11 / test " + "(default, 1, 14, runner)", + 1, + ) + ] + jobs.extend( + job( + f"linux-jammy-cuda13.2-py3.10-gcc11 / test " + f"(default, {shard}, 14, runner)", + 100 + shard, + ) + for shard in range(1, 15) + ) + + family = choose_test_job_family( + jobs, + "default", + "cuda", + "linux-jammy-cuda13.0-py3.10-gcc11", + ) + + self.assertEqual(family["prefix"], "linux-jammy-cuda13.2-py3.10-gcc11") + + def test_discovers_historical_rocm_family(self): + jobs = [ + job( + f"linux-jammy-rocm-py3.10-mi350 / test " + f"(distributed, {shard}, 3, runner)", + shard, + ) + for shard in range(1, 4) + ] + + family = choose_test_job_family( + jobs, + "distributed", + "rocm", + "linux-noble-rocm-py3.11-mi350", + ) + + self.assertEqual(family["prefix"], "linux-jammy-rocm-py3.10-mi350") + self.assertEqual(family["total"], 3) + + def test_preserves_discovered_test_kind(self): + jobs = [ + job( + f"linux-jammy-cuda13.2-py3.10-gcc11 / test-osdc " + f"(distributed, {shard}, 10, runner)", + shard, + ) + for shard in range(1, 11) + ] + + family = choose_test_job_family( + jobs, + "distributed", + "cuda", + "linux-jammy-cuda13.2-py3.10-gcc11", + ) + + self.assertEqual(family["kind"], "test-osdc") + + def test_marks_missing_shards_incomplete(self): + jobs = [ + job( + "linux-jammy-cuda13.2-py3.10-gcc11 / test " + "(default, 1, 14, runner)", + 1, + ) + ] + + family = choose_test_job_family( + jobs, + "default", + "cuda", + "linux-jammy-cuda13.2-py3.10-gcc11", + ) + + self.assertFalse(family["complete"]) + + def test_keeps_platforms_separate(self): + jobs = [ + job( + "linux-jammy-cuda13.2-py3.10-gcc11 / test " + "(default, 1, 14, runner)", + 1, + ) + ] + + self.assertIsNone( + choose_test_job_family( + jobs, + "default", + "rocm", + "linux-noble-rocm-py3.11-mi350", + ) + ) + + +if __name__ == "__main__": + unittest.main() From c4ea3db42ef382c55453ed5eefde600d36dc170b Mon Sep 17 00:00:00 2001 From: ethanwee1 Date: Fri, 18 Sep 2026 15:31:22 +0000 Subject: [PATCH 3/3] [CI] Scope parity auto-trigger to trunk runs Carry the exact trunk run ID into readiness checks and use a version-independent CUDA family matcher so periodic CUDA 13.4 jobs cannot satisfy a trunk CUDA 13.2 gate. --- .../parity_job_config.json | 2 +- .github/workflows/parity-auto.yml | 22 +++++++++---------- 2 files changed, 12 insertions(+), 12 deletions(-) diff --git a/.automation_scripts/pytorch-unit-test-scripts/parity_job_config.json b/.automation_scripts/pytorch-unit-test-scripts/parity_job_config.json index 6c6c38254e6ca..aa4a220095fae 100644 --- a/.automation_scripts/pytorch-unit-test-scripts/parity_job_config.json +++ b/.automation_scripts/pytorch-unit-test-scripts/parity_job_config.json @@ -16,7 +16,7 @@ "inductor": [{ "workflow": "inductor", "job_prefix": "unit-test / inductor-test" }], "test_kinds": ["test-osdc", "test"], "shard_counts": { "default": 14, "distributed": 10, "inductor": 2 }, - "checkrun_regex": "(linux-jammy-cuda13[.]2-py3[.]10-gcc11 / (test-osdc|test) [(](default|distributed),|unit-test / inductor-test / (test-osdc|test) [(]inductor,)" + "checkrun_regex": "(^linux-[^ ]*-cuda[0-9]+[.][0-9]+-py3[.][0-9]+-gcc[0-9]+ / (test-osdc|test) [(](default|distributed),|unit-test / inductor-test / (test-osdc|test) [(]inductor,)" }, "rocm": { "mi350": { diff --git a/.github/workflows/parity-auto.yml b/.github/workflows/parity-auto.yml index 06eb450f9bdbf..54931d01bc4f0 100644 --- a/.github/workflows/parity-auto.yml +++ b/.github/workflows/parity-auto.yml @@ -134,7 +134,7 @@ jobs: "" } - # Echo " " lines for the most recent completed + # Echo " " lines for recent completed # trunk.yml pushes, deduped and newest-first, up to MAX_COMMITS. # Candidate SHAs come from completed trunk pushes (not raw main # commits): the report consumes trunk's jobs, so a completed trunk run @@ -144,7 +144,7 @@ jobs: while [ "$(echo "$commits_json" | jq 'length')" -lt "$MAX_COMMITS" ]; do page_runs=$(gh api \ "repos/$UPSTREAM/actions/workflows/trunk.yml/runs?branch=$BRANCH&event=push&status=completed&per_page=100&page=$page" \ - --jq '.workflow_runs | map({head_sha, created_at})' 2>/dev/null) + --jq '.workflow_runs | map({id, head_sha, created_at})' 2>/dev/null) # On a gh api failure (rate limit / transient 5xx) page_runs is # empty or non-JSON; stop paginating and use what we have so far. if ! echo "$page_runs" | jq -e . >/dev/null 2>&1; then @@ -163,7 +163,7 @@ jobs: ' <(echo "$commits_json") <(echo "$page_runs")) page=$((page + 1)) done - echo "$commits_json" | jq -r '.[] | "\(.head_sha) \(.created_at)"' + echo "$commits_json" | jq -r '.[] | "\(.id) \(.head_sha) \(.created_at)"' } # Echo a JSON array of recent auto-parity workflow_dispatch runs in our @@ -250,7 +250,7 @@ jobs: # not already processed. Returns 0 to keep scanning, 1 to stop the scan # (commit too old, or per-scan dispatch cap reached). process_commit() { - local sha="$1" date="$2" + local run_id="$1" sha="$2" date="$3" local short commit_epoch all_check_runs cuda_check_runs local total_cr pending_cr pending_sample arch_dispatch short=$(echo "$sha" | cut -c1-8) @@ -273,11 +273,11 @@ jobs: return 0 fi - # Per-shard check-run state for this SHA (workflow_run conclusion can - # flip to failure before sibling shards finish, so we look per shard). + # Inspect only this trunk push. Commit-wide checks can include a + # periodic CUDA family with the same SHA and a different version. all_check_runs=$(gh api --paginate \ - "repos/$UPSTREAM/commits/$sha/check-runs?per_page=100" \ - --jq '.check_runs[] | {name,status,conclusion}' \ + "repos/$UPSTREAM/actions/runs/$run_id/jobs?filter=all&per_page=100" \ + --jq '.jobs[] | {name,status,conclusion}' \ 2>/dev/null | jq -s '.' || echo '[]') ROCM_CHECK_RUNS='[]' @@ -399,11 +399,11 @@ jobs: DISPATCHED_COUNT=0 DISPATCHED_SUMMARY="" - local sha date - while IFS=' ' read -r sha date; do + local run_id sha date + while IFS=' ' read -r run_id sha date; do [ -z "$sha" ] && continue # process_commit returns non-zero to stop the scan (too old / cap). - process_commit "$sha" "$date" || break + process_commit "$run_id" "$sha" "$date" || break done <<< "$commits" write_step_summary