Coverage for simulator/provisioning.py: 84%

165 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-09 04:47 +0000

1""" 

2Provisioning simulation module. 

3""" 

4from __future__ import annotations 

5 

6from tqdm.auto import tqdm 

7 

8import logging 

9 

10from typing import Optional 

11 

12from itertools import product 

13from itertools import combinations_with_replacement 

14 

15from functools import partial 

16 

17from concurrent.futures import ProcessPoolExecutor 

18from concurrent.futures import TimeoutError 

19from concurrent.futures import as_completed 

20 

21from sim_types import WorkflowConfig 

22from sim_types import GPUType 

23from sim_types import LatencyData 

24from sim_types import Provision 

25from sim_types import ProvisioningResult 

26from sim_types import Model 

27from sim_types import ModelAllocation 

28from sim_types import PowerData 

29from sim_types import QualityLevel 

30from sim_types import Policy 

31from sim_types import Result 

32from sim_types import num_gpus_to_str 

33 

34from auto_model_allocator import AutoModelAllocator 

35 

36from model_provisioner.policies import STREAMWISE_POLICY 

37 

38from constants import SECONDS_IN_HOUR 

39 

40 

41GPU_PROVISIONS: list[int] = [ 

42 0, 

43 8, 16, 24, 32, 40, 48, 56, 64, 72, 80, 88, 

44 96, 104, 112, 128, 144, 160, 176, 192, 208, 

45 224, 240, 256, 288, 304, 320, 352, 384, 400, 448, 480, 

46 512, 576, 624, 640, 672, 704, 768, 832, 864, 896, 960, 

47 1024, 1152, 1280, 1408, 1536, 1664, 1792, 2048, 2304, 2560, 2880, 

48 3200, 3584, 4096, 

49] 

50 

51 

52# Trimmed down 

53GPU_PROVISIONS_SHORT: list[int] = [ 

54 0, 

55 8, 16, 32, 48, 64, 72, 

56 96, 128, 192, 

57 256, 480, 512, 768, 1024, 2048, 4096, 

58] 

59 

60 

61def get_provisions( 

62 gpus_types: list[GPUType], 

63 limits_pairs: bool = True, 

64 short_list: bool = False, 

65) -> list[Provision]: 

66 """ 

67 Generate a list of provisioning options for the given GPU types. 

68 If limits_pairs=True, we only support pairs. 

69 """ 

70 assert 0 < len(gpus_types) 

71 

72 if len(gpus_types) <= 2 or not limits_pairs: 

73 return get_provisions_internal( 

74 gpus_types, 

75 short_list=short_list) 

76 

77 # Generate all the pairs of GPU types 

78 provisions: list[Provision] = [] 

79 for gpu_type_pair in combinations_with_replacement(gpus_types, r=2): 

80 assert len(gpu_type_pair) == 2 

81 if gpu_type_pair[0] == gpu_type_pair[1]: 

82 continue # single GPU type 

83 pair_provisions = get_provisions_internal( 

84 list(gpu_type_pair), 

85 short_list=short_list) 

86 provisions.extend(pair_provisions) 

87 provisions = remove_duplicate_provisions(provisions) 

88 return provisions 

89 

90 

91def get_provisions_internal( 

92 gpus_types: list[GPUType], 

93 short_list: bool = False, 

94) -> list[Provision]: 

95 """ 

96 Generate a list of provisioning options for the given GPU types. 

97 """ 

98 provisions: list[Provision] = [] 

99 

100 for counts in product(GPU_PROVISIONS_SHORT if short_list else GPU_PROVISIONS, repeat=len(gpus_types)): 

101 num_gpus = { 

102 gpu_type: count 

103 for gpu_type, count in zip(gpus_types, counts) 

104 if count > 0 

105 } 

106 if len(num_gpus) > 0: 

107 provisions.append(Provision(num_gpus=num_gpus)) 

108 

109 provisions = remove_duplicate_provisions(provisions) 

110 

111 return provisions 

112 

113 

114def remove_duplicate_provisions( 

115 provisions: list[Provision], 

116) -> list[Provision]: 

117 unique_provisions_dict: dict[str, Provision] = {} 

118 for provision in provisions: 

119 key = ','.join([ 

120 f"{gpu_type.value}:{count}" 

121 for gpu_type, count in sorted(provision.num_gpus.items()) 

122 ]) 

123 unique_provisions_dict[key] = provision 

124 

125 return list(unique_provisions_dict.values()) 

126 

127 

128def _process_provision( 

129 provision: Provision, 

130 workflow: WorkflowConfig, 

131 latency_data: LatencyData, 

132 power_data: Optional[PowerData], 

133 policy: Policy, 

134 verbose: bool, 

135) -> tuple[Provision, Optional[Result]]: 

136 """Worker function to process a single provision.""" 

137 gpu_types = [ 

138 gpu_type 

139 for gpu_type, count in provision.num_gpus.items() 

140 if count > 0 

141 ] 

142 assert 0 < len(gpu_types), f"No GPUs provisioned: {provision.num_gpus}." 

143 assert len(gpu_types) <= 2, f"Only support up to 2 GPU types in a provision: {provision.num_gpus}." 

144 

145 try: 

146 allocator = AutoModelAllocator( 

147 workflow=workflow, 

148 latency_data=latency_data, 

149 power_data=power_data, 

150 policy=policy, 

151 ) 

152 result = allocator.allocate( 

153 num_gpus=provision.num_gpus, 

154 verbose=verbose, 

155 ) 

156 if verbose: 

157 logging.info( 

158 f"Total time for {provision}: " 

159 f"{result.total_time_s:.2f} seconds ({result.total_time_s / SECONDS_IN_HOUR:.2f} hours)") 

160 return (provision, result) 

161 except KeyError as key_ex: 

162 logging.error(f"Error processing provision {provision}: {key_ex}", exc_info=True) 

163 return (provision, None) 

164 

165 

166def get_provisioning_results( 

167 workflow: WorkflowConfig, 

168 latency_data: LatencyData, 

169 power_data: Optional[PowerData] = None, 

170 policy: Policy = STREAMWISE_POLICY, 

171 provisions: Optional[list[Provision]] = None, 

172 verbose: bool = False, 

173 max_workers: Optional[int] = None, 

174 timeout: float = 10.0, # 10 seconds max per provision 

175 short_list: bool = False, 

176) -> ProvisioningResult: 

177 """ 

178 Get provisioning results for a list of GPU options. 

179 

180 Args: 

181 max_workers: Maximum number of worker processes. None uses all available CPUs. 

182 timeout: Timeout in seconds for each provision task (default: 10.0). 

183 """ 

184 times: list[float] = [] 

185 costs: list[float] = [] 

186 energies: list[float] = [] 

187 ttffs: list[float] = [] 

188 tbfs: list[float] = [] 

189 

190 actual_provision: list[dict[GPUType, int]] = [] 

191 config_provision: list[dict[GPUType, int]] = [] 

192 model_provision: list[dict[GPUType, dict[Model, list[ModelAllocation]]]] = [] 

193 

194 if provisions is None: 

195 provisions = get_provisions( 

196 policy.hardware, 

197 short_list=short_list) 

198 

199 worker_func = partial( 

200 _process_provision, 

201 workflow=workflow, 

202 latency_data=latency_data, 

203 power_data=power_data, 

204 policy=policy, 

205 verbose=verbose, 

206 ) 

207 

208 with ProcessPoolExecutor(max_workers=max_workers) as executor: 

209 futures = { 

210 executor.submit(worker_func, provision): provision 

211 for provision in provisions 

212 } 

213 with tqdm(total=len(provisions), desc=f"{policy.name} policy") as pbar: 

214 for future in as_completed(futures): 

215 provision = futures[future] 

216 try: 

217 provision_result, result = future.result(timeout=timeout) 

218 if result is None: 

219 logging.error(f"Skipping provision {provision} due to errors.") 

220 else: 

221 times.append(result.total_time_s) 

222 costs.append(result.cost) 

223 energies.append(result.total_energy) 

224 ttffs.append(result.ttff_s) 

225 tbfs.append(result.tbf_s) 

226 actual_provision.append(result.gpus_used.copy()) 

227 config_provision.append(provision_result.num_gpus.copy()) 

228 model_provision.append(result.models) 

229 except TimeoutError: 

230 logging.warning(f"Provision {provision} timed out after {timeout} seconds.") 

231 finally: 

232 pbar.update(1) 

233 

234 # Sort results by the order of input provisions 

235 provision_to_index = {num_gpus_to_str(p.num_gpus): i for i, p in enumerate(provisions)} 

236 sorted_indices = sorted( 

237 range(len(config_provision)), 

238 key=lambda i: provision_to_index.get(num_gpus_to_str(config_provision[i]), float('inf')) 

239 ) 

240 times = [times[i] for i in sorted_indices] 

241 costs = [costs[i] for i in sorted_indices] 

242 energies = [energies[i] for i in sorted_indices] 

243 ttffs = [ttffs[i] for i in sorted_indices] 

244 tbfs = [tbfs[i] for i in sorted_indices] 

245 actual_provision = [actual_provision[i] for i in sorted_indices] 

246 config_provision = [config_provision[i] for i in sorted_indices] 

247 model_provision = [model_provision[i] for i in sorted_indices] 

248 

249 return ProvisioningResult( 

250 latencies=times, 

251 costs=costs, 

252 energies=energies, 

253 ttffs=ttffs, 

254 tbfs=tbfs, 

255 actual_provision=actual_provision, 

256 config_provision=config_provision, 

257 model_provision=model_provision, 

258 ) 

259 

260 

261def get_provisioning_adaptive_results( 

262 workflow_config: WorkflowConfig, 

263 provisioning_qualities: dict[QualityLevel, ProvisioningResult], 

264 video_seconds: int = 10 * 60, 

265) -> ProvisioningResult: 

266 """ 

267 Get provisioning results for adaptive quality policy. 

268 """ 

269 assert len(provisioning_qualities) == 3 

270 

271 num_provisions = len(provisioning_qualities[QualityLevel.HIGH].costs) 

272 

273 assert num_provisions == len(provisioning_qualities[QualityLevel.HIGH].actual_provision), \ 

274 f"High: {num_provisions} != {len(provisioning_qualities[QualityLevel.HIGH].actual_provision)}" 

275 assert num_provisions == len(provisioning_qualities[QualityLevel.MEDIUM].actual_provision), \ 

276 f"Medium: {num_provisions} != {len(provisioning_qualities[QualityLevel.MEDIUM].actual_provision)}" 

277 assert num_provisions == len(provisioning_qualities[QualityLevel.LOW].actual_provision), \ 

278 f"Low: {num_provisions} != {len(provisioning_qualities[QualityLevel.LOW].actual_provision)}" 

279 

280 total_frames_ft = workflow_config.total_frames[Model.FT] 

281 

282 # for each provisioning option, get the adaptive policy cost and TTFF 

283 costs: list[float] = [] 

284 energies: list[float] = [] 

285 latencies: list[float] = [] 

286 ttffs: list[float] = [] 

287 tbfs: list[float] = [] 

288 qualities: list[float] = [] 

289 

290 actual_provision: list[dict[GPUType, int]] = [] 

291 config_provision: list[dict[GPUType, int]] = [] 

292 model_provision: list[dict[GPUType, dict[Model, list[ModelAllocation]]]] = [] 

293 

294 for idx in range(num_provisions): 

295 num_gpus = provisioning_qualities[QualityLevel.HIGH].config_provision[idx] 

296 # Initial check 

297 gpu_types = [ 

298 gpu_type 

299 for gpu_type, count in num_gpus.items() 

300 if count > 0 

301 ] 

302 assert 0 < len(gpu_types), f"No GPUs provisioned: {num_gpus}." 

303 assert len(gpu_types) <= 2, f"Only support up to 2 GPU types in a provision: {num_gpus}." 

304 

305 config_provision.append(num_gpus.copy()) 

306 

307 # check the actual provision 

308 gpus_high = provisioning_qualities[QualityLevel.HIGH].actual_provision[idx] 

309 gpus_medium = provisioning_qualities[QualityLevel.MEDIUM].actual_provision[idx] 

310 gpus_low = provisioning_qualities[QualityLevel.LOW].actual_provision[idx] 

311 

312 # check if the TTFF of low is less than total time of video 

313 ttff_low = provisioning_qualities[QualityLevel.LOW].ttffs[idx] 

314 ttff_med = provisioning_qualities[QualityLevel.MEDIUM].ttffs[idx] 

315 if ttff_low > video_seconds: 

316 logging.warning( 

317 f"Cannot apply policy for {num_gpus_to_str(num_gpus)}. " 

318 f"TTFF low {ttff_low:.2f} > Video {video_seconds}.") 

319 raise ValueError(f"Low quality TTFF ({ttff_low:.2f}) exceeds video length ({video_seconds}).") 

320 

321 # the portion of the low, medium, and high quality video 

322 portion_low = int(ttff_low / video_seconds * total_frames_ft) 

323 portion_medium = int(ttff_med / video_seconds * total_frames_ft) 

324 if portion_medium > total_frames_ft: 

325 portion_medium = total_frames_ft - portion_low 

326 portion_high = total_frames_ft - portion_low - portion_medium 

327 

328 logging.debug( 

329 f"Adaptive policy for {num_gpus_to_str(num_gpus)}, " 

330 f"Portions Low:{portion_low} Medium:{portion_medium} High:{portion_high}") 

331 

332 # Options to calculate the adaptive cost: 

333 # 1. Calculate the adaptive cost proportionally (not used) 

334 adaptive_cost = ( 

335 portion_low / total_frames_ft * provisioning_qualities[QualityLevel.LOW].costs[idx] 

336 + portion_medium / total_frames_ft * provisioning_qualities[QualityLevel.MEDIUM].costs[idx] 

337 + portion_high / total_frames_ft * provisioning_qualities[QualityLevel.HIGH].costs[idx] 

338 ) 

339 # 2. Sum of all costs (not used) 

340 adaptive_cost = ( 

341 provisioning_qualities[QualityLevel.LOW].costs[idx] 

342 + provisioning_qualities[QualityLevel.MEDIUM].costs[idx] 

343 + provisioning_qualities[QualityLevel.HIGH].costs[idx] 

344 ) 

345 # 3. Calculate the adaptive cost with the most expensive setting (not used) 

346 if portion_high > 0: 

347 adaptive_cost = provisioning_qualities[QualityLevel.HIGH].costs[idx] 

348 elif portion_medium > 0: 

349 adaptive_cost = provisioning_qualities[QualityLevel.MEDIUM].costs[idx] 

350 else: 

351 adaptive_cost = provisioning_qualities[QualityLevel.LOW].costs[idx] 

352 # 4. Highest cost 

353 adaptive_cost = provisioning_qualities[QualityLevel.HIGH].costs[idx] 

354 

355 adaptive_energy = ( 

356 portion_low / total_frames_ft * provisioning_qualities[QualityLevel.LOW].energies[idx] 

357 + portion_medium / total_frames_ft * provisioning_qualities[QualityLevel.MEDIUM].energies[idx] 

358 + portion_high / total_frames_ft * provisioning_qualities[QualityLevel.HIGH].energies[idx] 

359 ) 

360 

361 costs.append(adaptive_cost) 

362 energies.append(adaptive_energy) 

363 latencies.append(provisioning_qualities[QualityLevel.LOW].latencies[idx]) 

364 ttffs.append(ttff_low) 

365 tbfs.append(provisioning_qualities[QualityLevel.LOW].tbfs[idx]) 

366 weighted_quality = ( 

367 portion_low / total_frames_ft * 1 

368 + portion_medium / total_frames_ft * 2 

369 + portion_high / total_frames_ft * 3 

370 ) 

371 qualities.append(weighted_quality) 

372 

373 actual_provision.append({ 

374 gpu_type: max( 

375 gpus_high.get(gpu_type, 0), 

376 gpus_medium.get(gpu_type, 0), 

377 gpus_low.get(gpu_type, 0)) 

378 for gpu_type in gpu_types 

379 }) 

380 model_provision.append(provisioning_qualities[QualityLevel.HIGH].model_provision[idx]) 

381 

382 return ProvisioningResult( 

383 latencies=latencies, 

384 costs=costs, 

385 energies=energies, 

386 ttffs=ttffs, 

387 tbfs=tbfs, 

388 actual_provision=actual_provision, 

389 config_provision=config_provision, 

390 qualities=qualities, 

391 model_provision=model_provision, 

392 )