-
Notifications
You must be signed in to change notification settings - Fork 50
Expand file tree
/
Copy pathengineV4.py
More file actions
5246 lines (4758 loc) · 204 KB
/
Copy pathengineV4.py
File metadata and controls
5246 lines (4758 loc) · 204 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
from __future__ import annotations
import argparse
import atexit
import gc
import heapq
import importlib
import itertools
import math
import multiprocessing as mp
import os
import queue
import re
import shlex
import shutil
import signal
import subprocess
import sys
import tempfile
import threading
import time
from collections import OrderedDict, deque
# ThreadPoolExecutor 仅用于给进程内 NVML 采样加超时边界。
from concurrent.futures import ThreadPoolExecutor
from concurrent.futures import TimeoutError as FuturesTimeoutError
from dataclasses import dataclass, field, replace
from datetime import datetime
from multiprocessing import cpu_count, set_start_method
from pathlib import Path
from types import SimpleNamespace
import numpy as np
import pynvml
import yaml
from sanitizer_session import (
SanitizerSession,
encode_ready,
encode_result,
parse_event,
)
from tester.reporting.dump_writer import (
dump_enabled,
parse_strict_bool,
record_dump_terminal_status,
resolve_dump_options,
)
from tester.runtime.gpu_memory_preflight import (
GpuMemoryDeferred,
estimate_gpu_memory,
should_check_grad,
)
GIB = 1024**3
# GPU 超时只由批次入口读取,避免单 case 和批次形成不同的环境变量优先级。
GPU_PRESSURE_TIMEOUT_ENV_VAR = "PADDLEAPITEST_GPU_PRESSURE_TIMEOUT_SECONDS"
# 600 秒是显存压力保护,不是单 case 的执行超时。
DEFAULT_GPU_PRESSURE_TIMEOUT_SECONDS = 600.0
# slot 尚未 ready 的初始化态:批处理主循环据此判断是否仍有进展中 worker。
_INITIALIZING_SLOT_STATES = frozenset({"starting", "loaded", "preparing"})
_SLOT_STARTABLE_STATES = frozenset({"dead", "suspended"})
_SLOT_SCHEDULABLE_STATES = frozenset({"idle", "dead", "suspended"})
_SERVICEABLE_SLOT_STATES = _INITIALIZING_SLOT_STATES | frozenset({"idle", "busy"})
# 初始化连续失败的退避阶梯;失败次数超出阶梯长度后重复使用最后一档作为封顶间隔。
WORKER_INIT_BACKOFF_SECONDS = (5.0, 15.0, 30.0, 60.0)
# 连续初始化失败达到上限的 slot 不再复活,避免单个故障 slot 反复吞掉同一个 case。
MAX_SLOT_INIT_FAILURES_ENV_VAR = "PADDLEAPITEST_MAX_SLOT_INIT_FAILURES"
DEFAULT_MAX_SLOT_INIT_FAILURES = 5
# ledger 二次确认失败触发的换血同样有上限:这条路径不经过 init_failures,
# 不加独立计数的话 slot 会以 serviceable 身份无限换血,是批次永挂的直接根因之一。
MAX_SLOT_STARTUP_CHURNS_ENV_VAR = "PADDLEAPITEST_MAX_SLOT_STARTUP_CHURNS"
DEFAULT_MAX_SLOT_STARTUP_CHURNS = 3
# 初始化超时的现场只存在于被杀之前,诊断开关默认开启以便偶发故障可事后归因。
INIT_TIMEOUT_DIAGNOSTICS_ENV_VAR = "PADDLEAPITEST_INIT_TIMEOUT_DIAGNOSTICS"
# 单次现场采集的总预算,避免按进程数累加阻塞 watchdog。
INIT_TIMEOUT_DIAGNOSTICS_BUDGET_SECONDS = 10.0
from tester.reporting import (
init_log,
log_aggregation,
log_report,
log_retest,
log_runtime,
log_worker,
)
from tester.runtime.config_file_loader import resolve_config_files
from tester.runtime.runtime_config import (
TestRuntimeConfig,
limit_worker_layout,
runtime_config_for_gpu,
)
from tester.runtime.sanitizer_output import analyze_sanitizer_output
os.environ["FLAGS_use_system_allocator"] = "1"
os.environ["NVIDIA_TF32_OVERRIDE"] = "0"
@dataclass(frozen=True)
class GpuMemorySnapshot:
# snapshot 是调度器唯一认可的物理显存输入,不保存进程级占用猜测。
"""一次物理显存采样;free 已包含外部进程和遗留 CUDA context。"""
total_bytes: int
free_bytes: int
@dataclass(frozen=True)
class CaseGpuEstimate:
# compute/comparison 分开记录,双卡模式不能用计算卡峰值代替对比卡预算。
compute_bytes: int = 0
comparison_bytes: int = 0
# worker 只需要首阶段非 plan 峰值做动态 headroom 检查,避免重复解析完整配置。
compute_headroom_bytes: int | None = None
@dataclass
class GpuReclaimTracker:
"""用两个独立物理采样确认异步释放已经停止。"""
tolerance_bytes: int = 64 * 1024**2
sample_interval_seconds: float = 1.0
_last_free_bytes: int | None = None
_last_sample_at: float | None = None
def reset(self):
# 新的 release_pending 周期必须重新建立物理显存基线。
self._last_free_bytes = None
self._last_sample_at = None
def observe(self, snapshot, *, now):
# NVML 快照可能在多个主循环迭代中重复,间隔不足时不算独立样本。
free_bytes = max(0, int(snapshot.free_bytes))
if self._last_sample_at is None:
self._last_free_bytes = free_bytes
self._last_sample_at = float(now)
return False
if float(now) - self._last_sample_at < self.sample_interval_seconds:
return False
stable = abs(free_bytes - int(self._last_free_bytes or 0)) <= self.tolerance_bytes
self._last_free_bytes = free_bytes
self._last_sample_at = float(now)
return stable
@dataclass
class GpuPressureTimeout:
"""记录 pending 在没有任何持久化进展时的连续阻塞时长。"""
timeout_seconds: float
blocked_since: float | None = None
def update(self, *, blocked, now):
# 任意真实派发或持久化终态都会清零连续阻塞窗口。
if not blocked:
self.blocked_since = None
return False
if self.blocked_since is None:
self.blocked_since = float(now)
return float(now) - self.blocked_since >= self.timeout_seconds
def read_gpu_pressure_timeout(environ=None):
# 非法环境变量在批次启动阶段失败,不能运行到压力循环才静默采用默认值。
source = os.environ if environ is None else environ
raw_value = source.get(
GPU_PRESSURE_TIMEOUT_ENV_VAR,
str(DEFAULT_GPU_PRESSURE_TIMEOUT_SECONDS),
)
try:
timeout = float(raw_value)
except (TypeError, ValueError) as err:
raise ValueError(
f"{GPU_PRESSURE_TIMEOUT_ENV_VAR} must be a finite non-negative number, "
f"got {raw_value!r}"
) from err
if not math.isfinite(timeout) or timeout < 0:
raise ValueError(
f"{GPU_PRESSURE_TIMEOUT_ENV_VAR} must be a finite non-negative number, "
f"got {raw_value!r}"
)
return timeout
def read_max_slot_init_failures(environ=None):
# 0 表示只退避、永不放弃 slot;非法值在批次启动阶段就失败,不留给调度循环。
source = os.environ if environ is None else environ
raw_value = source.get(
MAX_SLOT_INIT_FAILURES_ENV_VAR,
str(DEFAULT_MAX_SLOT_INIT_FAILURES),
)
try:
max_failures = int(raw_value)
except (TypeError, ValueError) as err:
raise ValueError(
f"{MAX_SLOT_INIT_FAILURES_ENV_VAR} must be a non-negative integer, got {raw_value!r}"
) from err
if max_failures < 0:
raise ValueError(
f"{MAX_SLOT_INIT_FAILURES_ENV_VAR} must be a non-negative integer, got {raw_value!r}"
)
return max_failures
def read_max_slot_startup_churns(environ=None):
# 0 表示只退避、永不因换血退役 slot;非法值在批次启动阶段就失败。
source = os.environ if environ is None else environ
raw_value = source.get(
MAX_SLOT_STARTUP_CHURNS_ENV_VAR,
str(DEFAULT_MAX_SLOT_STARTUP_CHURNS),
)
try:
max_churns = int(raw_value)
except (TypeError, ValueError) as err:
raise ValueError(
f"{MAX_SLOT_STARTUP_CHURNS_ENV_VAR} must be a non-negative integer, got {raw_value!r}"
) from err
if max_churns < 0:
raise ValueError(
f"{MAX_SLOT_STARTUP_CHURNS_ENV_VAR} must be a non-negative integer, got {raw_value!r}"
)
return max_churns
def read_no_progress_timeout(options, environ=None):
"""解析全局无进展超时;默认取 2 倍 case timeout 与 1800s 的较大者。"""
source = os.environ if environ is None else environ
raw_value = source.get(NO_PROGRESS_TIMEOUT_ENV_VAR)
if raw_value is None:
case_timeout = float(getattr(options, "timeout", 0) or 0)
return max(2.0 * case_timeout, DEFAULT_NO_PROGRESS_TIMEOUT_SECONDS)
try:
timeout = float(raw_value)
except (TypeError, ValueError) as err:
raise ValueError(
f"{NO_PROGRESS_TIMEOUT_ENV_VAR} must be a finite non-negative number, got {raw_value!r}"
) from err
if not math.isfinite(timeout) or timeout < 0:
raise ValueError(
f"{NO_PROGRESS_TIMEOUT_ENV_VAR} must be a finite non-negative number, got {raw_value!r}"
)
return timeout
def slot_init_backoff_seconds(init_failures):
"""把连续失败次数映射到退避间隔,超出阶梯长度后停在封顶值。"""
if init_failures <= 0:
return 0.0
index = min(init_failures, len(WORKER_INIT_BACKOFF_SECONDS)) - 1
return WORKER_INIT_BACKOFF_SECONDS[index]
def parse_sanitizer_timing_file(path):
# 观测文件损坏只丢弃样本,不能改变 sanitizer 的业务结果。
values = {}
try:
with open(path, encoding="utf-8") as timing_file:
for line in timing_file:
# 只读取制表符分隔的 phase,普通输出不会进入统计。
phase, separator, raw_duration = line.rstrip("\n").partition("\t")
if not separator:
continue
try:
duration = float(raw_duration)
except ValueError:
continue
if duration >= 0:
values[phase] = values.get(phase, 0.0) + duration
except OSError:
return {}
return values
@dataclass(frozen=True)
class GpuSchedulingPolicy:
safety_reserve_bytes_min: int = 2 * GIB
safety_reserve_fraction: float = 0.05
minimum_case_bytes: int = 1 * GIB
case_margin_bytes: int = 512 * 1024**2
case_multiplier: float = 1.25
def safety_reserve_bytes(self, snapshot):
# 安全余量同时覆盖小卡固定开销和大卡按比例增长的 workspace。
return max(
self.safety_reserve_bytes_min,
int(max(0, snapshot.total_bytes) * self.safety_reserve_fraction),
)
def case_admission_bytes(self, estimated_peak_bytes):
# 未知配置仍需占用最小预算,估算值较大时再叠加边际和倍率。
estimate = max(0, int(estimated_peak_bytes))
return max(
self.minimum_case_bytes,
estimate + self.case_margin_bytes,
int(estimate * self.case_multiplier),
)
@dataclass
class GpuReservation:
# state 保留 active 与 release_pending 的差异,后者仍然占用承诺。
assignment_id: int
slot_index: int
device_bytes: dict[int, int]
state: str = "active"
class GpuReservationLedger:
"""按 assignment 持有显存承诺,release_pending 期间禁止复用。"""
def __init__(self, device_ids, *, snapshots, max_workers, policy=None):
# 一个 group 要么是一张计算卡,要么是不可拆分的双卡 pair。
if not device_ids or len(device_ids) > 2:
raise ValueError("a reservation group must contain one or two devices")
self.device_ids = tuple(device_ids)
self.max_workers = max(0, int(max_workers))
self.policy = policy or GpuSchedulingPolicy()
self._reservations = {}
self._committed = dict.fromkeys(self.device_ids, 0)
self._trackers = {gpu_id: GpuReclaimTracker() for gpu_id in self.device_ids}
self._pending_release = set()
self._claimed = set()
# 从单 worker 起步,稳定完成后再逐步扩容,避免启动阶段瞬间压满显存。
self.target_workers = min(self.max_workers, 1)
self._success_streak = 0
# 初始化快照后才允许 reserve,避免无采样状态被当成无限容量。
self.update_snapshots(snapshots)
def update_snapshots(self, snapshots):
# 只保留本 group 的设备,双卡确认必须使用同一轮调用传入的快照。
self._snapshots = {gpu_id: snapshots[gpu_id] for gpu_id in self.device_ids}
def _available(self, gpu_id):
# free 已含外部占用和遗留 context,不再按进程归属做推断。
snapshot = self._snapshots[gpu_id]
return max(0, snapshot.free_bytes - self.policy.safety_reserve_bytes(snapshot))
def _requested(self, estimates):
# 双卡 reservation 同时扣计算卡和 comparison 卡的预算。
requested = {self.device_ids[0]: self.policy.case_admission_bytes(estimates.compute_bytes)}
if len(self.device_ids) == 2:
requested[self.device_ids[1]] = self.policy.case_admission_bytes(
estimates.comparison_bytes
)
return requested
def reserve(self, *, assignment_id, slot_index, estimates):
# 常驻 worker 的旧 lease 可以被同一 slot 的新 assignment 原子替换。
if assignment_id in self._reservations:
return None
requested = self._requested(estimates)
old_lease = next(
(
reservation
for reservation in self._reservations.values()
if reservation.slot_index == slot_index and reservation.state == "release_pending"
),
None,
)
active_count = sum(
reservation.state == "active" for reservation in self._reservations.values()
)
if active_count >= self.target_workers:
return None
committed = dict(self._committed)
if old_lease is not None:
for gpu_id, amount in old_lease.device_bytes.items():
committed[gpu_id] -= amount
if any(
committed[gpu_id] + requested[gpu_id] > self._available(gpu_id)
for gpu_id in self.device_ids
):
return None
if old_lease is not None:
self._pending_release.discard(old_lease.assignment_id)
self._reservations.pop(old_lease.assignment_id, None)
self._claimed.discard(old_lease.assignment_id)
self._committed.update(committed)
# pending lease 被同 slot 接管时,旧回收样本不能污染下一次 reclaim 基线。
for tracker in self._trackers.values():
tracker.reset()
reservation = GpuReservation(assignment_id, slot_index, requested)
self._reservations[assignment_id] = reservation
for gpu_id, amount in requested.items():
self._committed[gpu_id] += amount
return reservation
def record_result(self, msg_type):
"""按结果调整 group 并发上限;异常只收缩,且受用户硬上限约束。"""
if msg_type == "done":
self._success_streak += 1
if self._success_streak >= 2 and self.target_workers < self.max_workers:
self.target_workers += 1
self._success_streak = 0
return
if msg_type in {"deferred", "timeout", "crashed"}:
self._success_streak = 0
self.target_workers = max(1, min(self.target_workers, self.max_workers) - 1)
def confirm(self, assignment_id, snapshots):
# worker bootstrap 可能改变 free,派发前必须做一次独立物理确认。
reservation = self._reservations.get(assignment_id)
if reservation is None or reservation.state != "active":
return False
self.update_snapshots(snapshots)
# 确认检查整个 group 的承诺,避免多个延迟启动的 worker 分别通过后合计超卖。
if any(self._committed[gpu_id] > self._available(gpu_id) for gpu_id in self.device_ids):
self.mark_release_pending(assignment_id)
return False
return True
def mark_release_pending(self, assignment_id):
# 终态只开启回收观察,不直接释放显存承诺。
reservation = self._reservations.get(assignment_id)
if reservation is None:
return False
if reservation.state == "active":
reservation.state = "release_pending"
self._pending_release.add(assignment_id)
return True
def advance_reclaim(self, snapshots, *, now):
# 同组所有设备稳定后批量释放,避免双卡只回收一半就重新派发。
if not self._pending_release:
return 0
self.update_snapshots(snapshots)
stable = all(
self._trackers[gpu_id].observe(self._snapshots[gpu_id], now=now)
for gpu_id in self.device_ids
)
if not stable:
# 任一设备仍在变化,保留全部 pending reservation。
return 0
released = []
for assignment_id in tuple(self._pending_release):
reservation = self._reservations.pop(assignment_id, None)
self._pending_release.discard(assignment_id)
if reservation is None:
continue
for gpu_id, amount in reservation.device_bytes.items():
self._committed[gpu_id] -= amount
released.append(assignment_id)
for tracker in self._trackers.values():
# 下一批回收必须使用新的基线,不能复用上一批的稳定样本。
tracker.reset()
return tuple(released)
def claim_terminal(self, assignment_id):
# claimed 集合防止 timeout/crash 与 worker 正常终态重复结算。
reservation = self._reservations.get(assignment_id)
if reservation is None or assignment_id in self._claimed:
return False
self._claimed.add(assignment_id)
self.mark_release_pending(assignment_id)
return True
# 运行时透传给 test class 的选项白名单。
VALID_TEST_ARGS = {
"test_amp",
"test_backward",
"atol",
"rtol",
"accuracy_manual_threshold_config",
"record_accuracy_tolerance",
"operation_mode",
"bos_path",
"random_seed",
"bos_conf_path",
"bcecmd_path",
"bitwise_alignment",
"use_gpu_mode",
}
SANITIZER_FORWARD_ARGS = {
"accuracy",
"paddle_only",
"paddle_cinn",
"paddle_gpu_performance",
"torch_gpu_performance",
"paddle_torch_gpu_performance",
"accuracy_stable",
"accuracy_dual_gpu",
"accuracy_stable_dual_gpu",
"paddle_custom_device",
"custom_device_vs_gpu",
"custom_device_vs_gpu_mode",
"test_amp",
"test_cpu",
"use_cached_numpy",
"use_gpu_mode",
"atol",
"rtol",
"accuracy_manual_threshold_config",
"record_accuracy_tolerance",
"test_backward",
"show_runtime_status",
"random_seed",
"bitwise_alignment",
}
SANITIZER_FORWARD_ARGS_SORTED = tuple(sorted(SANITIZER_FORWARD_ARGS))
# 运行时错误标记,避免在每个 case 里重复构造。
OOM_ERROR_MARKERS = (
"cuda out of memory",
"out of memory error",
"resourceexhaustederror",
"out of memory",
"outofmemoryerror",
"cannot allocate memory",
"std::bad_alloc",
"bad allocation",
"memoryerror",
"cublas_status_alloc_failed",
)
CUDA_ERROR_MARKERS = (
"cuda error",
"memory corruption",
"illegal memory access",
"invalid configuration argument",
"invalid resource handle",
)
GPU_PERFORMANCE_MODES = (
"paddle_gpu_performance",
"torch_gpu_performance",
"paddle_torch_gpu_performance",
)
# 主模式互斥校验只看这些开关;dual 标志会先展开成对应主模式。
PRIMARY_TEST_MODES = (
"paddle_only",
"paddle_cinn",
"accuracy",
"paddle_gpu_performance",
"torch_gpu_performance",
"paddle_torch_gpu_performance",
"accuracy_stable",
"paddle_custom_device",
"custom_device_vs_gpu",
)
TORCH_REFERENCE_MODES = (
"accuracy",
"accuracy_stable",
"accuracy_dual_gpu",
"accuracy_stable_dual_gpu",
"torch_gpu_performance",
"paddle_torch_gpu_performance",
)
TORCH_UTILITY_MODES = TORCH_REFERENCE_MODES + (
"paddle_cinn",
"paddle_gpu_performance",
"paddle_custom_device",
"custom_device_vs_gpu",
)
GPU_MEMORY_PREFLIGHT_MODES = (
"accuracy_stable_dual_gpu",
"accuracy_dual_gpu",
"accuracy_stable",
"accuracy",
"paddle_only",
)
_BYTES_PER_GIB = 1024**3
# 选择测试类的优先级顺序。
TEST_CLASS_BY_OPTION = (
("paddle_only", "APITestPaddleOnly"),
("paddle_cinn", "APITestCINNVSDygraph"),
("accuracy_dual_gpu", "APITestAccuracy"),
("accuracy", "APITestAccuracy"),
("paddle_gpu_performance", "APITestPaddleGPUPerformance"),
("torch_gpu_performance", "APITestTorchGPUPerformance"),
("paddle_torch_gpu_performance", "APITestPaddleTorchGPUPerformance"),
("accuracy_stable_dual_gpu", "APITestAccuracyStable"),
("accuracy_stable", "APITestAccuracyStable"),
("paddle_custom_device", "APITestCustomDeviceVSCPU"),
("custom_device_vs_gpu", "APITestPaddleDeviceVSGPU"),
)
# 设备探测命令和缓存状态。
XPU_SMI_COMMAND = "xpu-smi"
XPU_SMI_DEVICE_PATTERN = r"^\|\s*(\d+)\s+\S"
ILUVATAR_SMI_COMMAND = "ixsmi"
ILUVATAR_SMI_DEVICE_PATTERN = r"^\|\s*(\d+)\s+Iluvatar"
DEVICE_TYPE = None
DEVICE_TYPE_DETECTED = False
DEVICE_COUNT = None # 设备总数
_MEM_SNAPSHOT = None # gpu_id -> (total_gb, used_gb)
_MEM_SNAPSHOT_TS = 0.0
_NVML_INITIALIZED = False # 重复显存查询的 NVML 会话。
_MEM_SNAPSHOT_TTL = 2.0 # 秒。
# 调度与重试上限。
MAX_TOTAL_WORKERS = 64
MAX_EXTERNAL_KILL_RETRIES_PER_CASE = 1
MAX_TOTAL_EXTERNAL_KILL_EVENTS = 3
# 初始 warmup 与单个 slot 复活共用的启动超时预算。
WORKER_STARTUP_TIMEOUT = 180
# 全局 no-progress watchdog:独立于主循环存活,专治主进程卡 NVML/join/select
# 时既有 watchdog 全盲的场景。0 表示关闭。
NO_PROGRESS_TIMEOUT_ENV_VAR = "PADDLEAPITEST_NO_PROGRESS_TIMEOUT_SECONDS"
DEFAULT_NO_PROGRESS_TIMEOUT_SECONDS = 1800.0
NO_PROGRESS_CHECK_INTERVAL_SECONDS = 10.0
# watchdog 触发后给主线程自然退出的宽限期,超时则强制 os._exit(2)。
NO_PROGRESS_EXIT_GRACE_SECONDS = 60.0
FORECAST_MIN_INTERVAL_SECONDS = 60
FORECAST_MAX_INTERVAL_SECONDS = 30 * 60
FORECAST_TARGET_CASES = 100
FORECAST_INITIAL_MAX_WAIT_SECONDS = 5 * 60
GPU_MEMORY_DEFER_INITIAL_BACKOFF_SECONDS = 1.0
GPU_MEMORY_DEFER_MAX_BACKOFF_SECONDS = 30.0
# 首次 deferred 只清理并复用常驻 worker,连续失败才退役进程。
GPU_MEMORY_DEFER_RETIRE_AFTER = 2
SANITIZER_COMPUTE_BUDGET_ENV = "PADDLEAPITEST_SANITIZER_COMPUTE_BUDGET_GIB"
SANITIZER_COMPARISON_BUDGET_ENV = "PADDLEAPITEST_SANITIZER_COMPARISON_BUDGET_GIB"
def _is_unavailable_gpu_error(error_msg):
"""识别设备级初始化失败;普通 case 异常仍沿用原有重试路径。"""
text = str(error_msg).lower()
return any(
marker in text
for marker in (
"cudaerrordevicesunavailable",
"cuda error(46)",
"cudaerrordevicelost",
"cuda error(45)",
)
)
SANITIZER_TIMING_FILE_ENV = "PADDLEAPITEST_SANITIZER_TIMING_FILE"
# 调度只需要有限候选;窗口随最大并发放大,保留跳过大 case 的能力。
CANDIDATE_WINDOW_PER_WORKER = 32
CANDIDATE_WINDOW_MIN = 64
# 窗口内全部装不下时放大一次扫描范围,但仍然有界:无界扫描在百万级 pending 上
# 单轮就要十几秒,比它要修的问题更慢。超出该上界的 case 靠队首消费自然前移。
EXTENDED_SCAN_MAX_CANDIDATES = 8192
# 扩展扫描的最小间隔,从上一次扫描“结束”开始计时,避免扫描本身把间隔耗尽。
EXTENDED_SCAN_MIN_INTERVAL_SECONDS = 5.0
WORKER_PREPARE_RUNTIME = "prepare_runtime"
def _option_enabled(options, name):
return bool(getattr(options, name, False))
def _memory_defer_delay(retry_count):
retry_count = max(0, int(retry_count))
return min(
GPU_MEMORY_DEFER_MAX_BACKOFF_SECONDS,
GPU_MEMORY_DEFER_INITIAL_BACKOFF_SECONDS * (2 ** min(retry_count, 5)),
)
def _collect_wait_timeout(pending_dispatch, default=0.5):
pending_delay = pending_dispatch.earliest_delay()
if pending_delay is None:
return default
return min(default, max(0.1, pending_delay))
def _fatal_log_type_for_error(error, terminal_log_type=None):
"""按普通运行路径的优先级识别 fatal 类型。"""
error_text = str(error).lower()
if any(marker in error_text for marker in OOM_ERROR_MARKERS):
return "oom"
if terminal_log_type == "torch_error" and any(
marker in error_text for marker in CUDA_ERROR_MARKERS
):
return "torch_error"
if any(marker in error_text for marker in CUDA_ERROR_MARKERS):
return "paddle_cuda"
return None
def _normalize_sanitizer_exitcode(output, *, returncode, sanitizer_error_exitcode):
"""把 sanitizer 专用退出码还原为普通 worker 的 fatal 退出协议。"""
if returncode != sanitizer_error_exitcode:
return returncode
# 仅把 sanitizer 的协议码转换为统一 fatal 码,普通 child 退出码保持原语义。
# Torch 已写出的终态会出现在 child 输出中;OOM 仍由共享文本规则优先识别。
terminal_log_type = "torch_error" if "[torch_error]" in output.lower() else None
fatal_log_type = _fatal_log_type_for_error(output, terminal_log_type)
return log_worker.fatal_exit_code(fatal_log_type or "paddle_cuda", False)
def _sanitizer_analysis_is_ignored(analysis):
"""仅在过滤后没有应用 fatal 证据时忽略 sanitizer 噪声。"""
return bool(
analysis.only_ignored_diagnostics and _fatal_log_type_for_error(analysis.output) is None
)
def _sanitizer_session_output_has_report(output):
"""识别 session 当前 request 是否包含 compute-sanitizer 诊断块。"""
# 普通 Paddle 输出不会使用 sanitizer 的固定分隔线;按 request 片段检查可恢复逐 case 归因。
if "========= " not in output:
return False
if "Error:" in output or "Program hit" in output:
return True
summary_counts = re.findall(r"ERROR SUMMARY:\s*(\d+)\s+errors?", output)
return any(int(count) > 0 for count in summary_counts)
@dataclass
class GpuEstimateFailureReport:
"""汇总显存估算失败。
估算失败会让 case 退化到 1 GiB 准入下界:既可能过量准入导致真实 OOM,也可能
在 worker 侧被误判为 skip。逐条打印会在大批量失败时刷屏,因此只保留首条诊断
加总数,足以区分“个别 API 没有模型”和“估算器整体坏了”。
"""
total: int = 0
first_config: str | None = None
first_error: str | None = None
error_counts: dict[str, int] = field(default_factory=dict)
def record(self, config, err):
self.total += 1
error_kind = type(err).__name__
self.error_counts[error_kind] = self.error_counts.get(error_kind, 0) + 1
if self.first_config is None:
self.first_config = config
self.first_error = f"{type(err).__name__}: {err}"
def emit(self, all_case):
if not self.total:
return
print(
f"[gpu] ESTIMATE_FALLBACK | fallback={self.total}/{all_case} | "
f"types={self.error_counts} | first={self.first_config} | "
f"error={self.first_error}",
flush=True,
)
@dataclass
class BatchRetryState:
per_case_external_kill_retries: dict[str, int] = field(default_factory=dict)
per_case_memory_defer_retries: dict[str, int] = field(default_factory=dict)
case_memory_estimates: dict[str, CaseGpuEstimate] = field(default_factory=dict)
slot_memory_defer_retries: dict[int, int] = field(default_factory=dict)
total_external_kills: int = 0
unsafe_environment: bool = False
@dataclass
class BatchRunState:
tested_case: int = 0
batch_exit_code: int = 0
shutdown_force: bool = False
abort_run: bool = False
active_tasks: int = 0
test_started_at: float | None = None
last_forecast_at: float | None = None
last_forecast_case: int = 0
# 累计派发数:与 tested_case 共同构成 no-progress watchdog 的 progress 定义,
# 单调递增,只由主循环写入、watchdog 只读。
total_dispatched: int = 0
@dataclass(frozen=True)
class PendingCase:
"""待重试 case 及其最早可再次派发时间。"""
config: str
ready_at: float = 0.0
compute_estimate_bytes: int = 0
comparison_estimate_bytes: int = 0
# None 表示调度进程未获得可信估算,worker 必须执行完整预检。
compute_headroom_bytes: int | None = None
@property
def gpu_estimate(self):
return CaseGpuEstimate(
self.compute_estimate_bytes,
self.comparison_estimate_bytes,
self.compute_headroom_bytes,
)
class PendingQueue:
"""待派发 case 队列:ready 段保持 FIFO,延迟重试段按 ready_at 组织为最小堆。
协议边界:
- 所有读取入口内部先 promote 到期项,调用方无需关心两段结构的同步顺序。
- appendleft 用于整波回滚与外部 kill 重试;仅当 ready_at 已到期才真正插到队首。
- 与旧的单一 deque 实现的差异:未到期 case 不再占据 ready 段,因此取任务和
计算最近可派发时间都不需要遍历整个队列。
"""
__slots__ = (
"_delayed",
"_ready",
"_sequence",
"_ready_version",
"_ready_snapshot_cache",
"_ready_snapshot_iterator",
"_ready_snapshot_version",
)
def __init__(self, cases=()):
# OrderedDict 同时提供稳定 FIFO 和按 config 的 O(1) 删除。
self._ready = OrderedDict()
self._delayed = []
# 单调序号让相同 ready_at 的堆序稳定,同时避免堆比较落到 PendingCase 上。
self._sequence = 0
self._ready_version = 0
self._ready_snapshot_cache = None
self._ready_snapshot_iterator = None
self._ready_snapshot_version = -1
# 构造阶段没有并发读者,批量填充后一次发布版本,避免逐项失效扫描缓存。
for case in cases:
if case.ready_at <= 0.0:
self._ready[case.config] = case
else:
self._push_delayed(case)
if self._ready:
self._ready_version = 1
def __len__(self):
return len(self._ready) + len(self._delayed)
def __bool__(self):
return bool(self._ready) or bool(self._delayed)
@property
def ready_count(self):
"""返回已到期段长度;扩展扫描不应把延迟 case 算入候选窗口。"""
self._promote(time.monotonic())
return len(self._ready)
def _invalidate_ready_snapshot(self):
self._ready_version += 1
self._ready_snapshot_cache = None
self._ready_snapshot_iterator = None
self._ready_snapshot_version = -1
def _ready_snapshot_to(self, end):
"""按需缓存扫描前缀,避免百万级 ready 队列首次扩展时整队复制。"""
if self._ready_snapshot_version != self._ready_version:
self._ready_snapshot_cache = []
self._ready_snapshot_iterator = iter(self._ready.values())
self._ready_snapshot_version = self._ready_version
cache = self._ready_snapshot_cache
if cache is None or self._ready_snapshot_iterator is None:
return []
missing = max(0, int(end) - len(cache))
if missing:
cache.extend(itertools.islice(self._ready_snapshot_iterator, missing))
return cache
def _push_delayed(self, case):
self._sequence += 1
heapq.heappush(self._delayed, (case.ready_at, self._sequence, case))
def _promote(self, now):
promoted = []
while self._delayed and self._delayed[0][0] <= now:
promoted.append(heapq.heappop(self._delayed)[2])
if promoted:
# 重试 case 到期后优先于尚未派发的新 case;逆序扩展才能保留到期顺序。
for case in reversed(promoted):
self._ready[case.config] = case
self._ready.move_to_end(case.config, last=False)
self._invalidate_ready_snapshot()
def append(self, case, *, now=None):
now = time.monotonic() if now is None else now
if case.ready_at <= now:
self._ready[case.config] = case
self._invalidate_ready_snapshot()
else:
self._push_delayed(case)
def appendleft(self, case, *, now=None):
now = time.monotonic() if now is None else now
if case.ready_at <= now:
self._ready[case.config] = case
self._ready.move_to_end(case.config, last=False)
self._invalidate_ready_snapshot()
else:
self._push_delayed(case)
def clear(self):
self._ready.clear()
self._delayed.clear()
self._invalidate_ready_snapshot()
def earliest_delay(self, now=None):
"""返回最近一次可派发的剩余等待秒数;队列为空返回 None。
读取路径会顺带把到期的延迟 case 迁入 ready 段。
"""
if not self:
return None
now = time.monotonic() if now is None else now
self._promote(now)
if self._ready:
return 0.0
return max(0.0, self._delayed[0][0] - now)
def pop_ready(self, now=None):
"""取出一个已到期 case 的配置字符串;没有已到期 case 时返回 None。"""
now = time.monotonic() if now is None else now
self._promote(now)
if not self._ready:
return None
_, case = self._ready.popitem(last=False)
self._invalidate_ready_snapshot()
return case.config
def candidate_window(self, limit, now=None):
"""返回队首至多 limit 个已到期 case;limit 为 None 表示全部 ready case。
返回的是队首快照副本,不消费队列;真正取走由 take_case_selection 完成。
读取路径会顺带把到期的延迟 case 迁入 ready 段。
"""
now = time.monotonic() if now is None else now
self._promote(now)
if limit is None or limit >= len(self._ready):
return list(self._ready.values())
return list(itertools.islice(self._ready.values(), limit))
def scan_window(self, limit, *, cursor=0, now=None):
"""从 ready 段的游标位置取有界窗口,并返回下一次扫描游标。
ready 队列不变时复用一次快照,避免为了跨过一批不可准入 case 而重复从队首
扫描;队列结构变化后会让调用方重新从当前窗口起点开始。
"""
now = time.monotonic() if now is None else now
self._promote(now)
ready_count = len(self._ready)
if not ready_count:
return [], 0
start = int(cursor) % ready_count
end = min(start + max(0, int(limit)), ready_count)
ready = self._ready_snapshot_to(end)
next_cursor = end if end < ready_count else 0
return ready[start:end], next_cursor
def take_case_selection(self, candidates, selected_indices):
"""从任意扫描窗口中取走选中的 case,并保持其余 ready 顺序。"""
selected = [candidates[index] for index in selected_indices if 0 <= index < len(candidates)]
if not selected:
return []
# 只按 key 删除选中项,延迟堆和未选中的 ready 顺序都不变。
for case in selected:
self._ready.pop(case.config, None)
self._invalidate_ready_snapshot()
return selected
def iter_all(self):
"""遍历全部 case,仅用于构建批次级快照。"""
yield from self._ready.values()
for _, _, case in self._delayed:
yield case
@dataclass
class TerminalClaims:
"""case 终态的一次性认领表,CPU 和 GPU 路径共用同一协议。"""
tokens: dict[int, tuple[int | None, str]] = field(default_factory=dict)
def register(self, slot_index, worker_pid, config):
self.tokens[slot_index] = (worker_pid, config)
def claim(self, msg):
message = BatchMessage.from_raw(msg)
expected_token = self.tokens.get(message.slot_index)
if expected_token is None or expected_token != (message.worker_pid, message.config):
return False
del self.tokens[message.slot_index]
return True
@dataclass(frozen=True)
class WorkerTask:
"""主调度器传给 worker 的 case 及本波实际显存预算。"""
config: str
workers_on_gpu: int
compute_budget_gib: float
comparison_budget_gib: float = 0.0
compute_estimate_bytes: int = 0
comparison_estimate_bytes: int = 0
# 只跨进程传紧凑峰值,避免序列化完整配置分析结果。
compute_headroom_bytes: int | None = None
@dataclass
class CaseRuntimeContext:
started_at: float
gpu_id: int
comparison_gpu_id: int | None
suppress_case_tags: bool
runtime_config: object | None = None
@dataclass
class BatchConfigLoadResult: