Coverage for simulator/provisioning.py: 84%
165 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-09 04:47 +0000
« 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
6from tqdm.auto import tqdm
8import logging
10from typing import Optional
12from itertools import product
13from itertools import combinations_with_replacement
15from functools import partial
17from concurrent.futures import ProcessPoolExecutor
18from concurrent.futures import TimeoutError
19from concurrent.futures import as_completed
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
34from auto_model_allocator import AutoModelAllocator
36from model_provisioner.policies import STREAMWISE_POLICY
38from constants import SECONDS_IN_HOUR
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]
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]
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)
72 if len(gpus_types) <= 2 or not limits_pairs:
73 return get_provisions_internal(
74 gpus_types,
75 short_list=short_list)
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
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] = []
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))
109 provisions = remove_duplicate_provisions(provisions)
111 return provisions
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
125 return list(unique_provisions_dict.values())
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}."
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)
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.
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] = []
190 actual_provision: list[dict[GPUType, int]] = []
191 config_provision: list[dict[GPUType, int]] = []
192 model_provision: list[dict[GPUType, dict[Model, list[ModelAllocation]]]] = []
194 if provisions is None:
195 provisions = get_provisions(
196 policy.hardware,
197 short_list=short_list)
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 )
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)
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]
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 )
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
271 num_provisions = len(provisioning_qualities[QualityLevel.HIGH].costs)
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)}"
280 total_frames_ft = workflow_config.total_frames[Model.FT]
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] = []
290 actual_provision: list[dict[GPUType, int]] = []
291 config_provision: list[dict[GPUType, int]] = []
292 model_provision: list[dict[GPUType, dict[Model, list[ModelAllocation]]]] = []
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}."
305 config_provision.append(num_gpus.copy())
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]
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}).")
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
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}")
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]
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 )
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)
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])
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 )