diff --git a/monitoring/benchmarker/engine/loads/step_execution.py b/monitoring/benchmarker/engine/loads/step_execution.py index 95c2b6c1e8..80c6c3dd44 100644 --- a/monitoring/benchmarker/engine/loads/step_execution.py +++ b/monitoring/benchmarker/engine/loads/step_execution.py @@ -3,7 +3,7 @@ from collections.abc import Callable from concurrent.futures import ThreadPoolExecutor from dataclasses import dataclass, field -from datetime import UTC, datetime +from datetime import UTC, datetime, timedelta from random import Random from typing import Any @@ -36,6 +36,9 @@ PERIODIC_STATUS_PERIOD_S = 30.0 +OPERATION_ORDER_TOLERANCE = timedelta(seconds=60) +"""Operations are recorded as they complete, so they are ordered by completion time to within this tolerance.""" + @dataclass class ActiveUser: @@ -198,6 +201,21 @@ async def wind_down_and_cleanup_remaining_users( ) +def first_index_completed_since( + operations: list[ExecutedOperation], t: datetime +) -> int: + """Return an index such that every operation completed at or after `t` is at or after that index in `operations`. + + Operations are appended to `operations` as they complete, so they are ordered by completion time to within + OPERATION_ORDER_TOLERANCE. Operations completed before `t` may still be present after the returned index. + """ + cutoff = t - OPERATION_ORDER_TOLERANCE + i = len(operations) + while i > 0 and operations[i - 1].completed_at.datetime >= cutoff: + i -= 1 + return i + + async def run_load_step( load: UserBasedLoad, load_label: str, @@ -216,6 +234,8 @@ async def run_load_step( current_virtual_users = [au.user for au in active_users] + step_ops_start = first_index_completed_since(operations, step_start_time) + # Wait for throughput to become stable for this step stability_time: datetime | None = None is_unstable = False @@ -227,7 +247,7 @@ async def run_load_step( and load.throughput_instability_criteria and check_stability_criteria( load.throughput_instability_criteria, - operations, + operations[step_ops_start:], current_virtual_users, step_start_time, now, @@ -242,7 +262,7 @@ async def run_load_step( if check_stability_criteria( load.throughput_stability_criteria, - operations, + operations[step_ops_start:], current_virtual_users, step_start_time, now, @@ -270,7 +290,7 @@ async def run_load_step( and load.throughput_instability_criteria and check_stability_criteria( load.throughput_instability_criteria, - operations, + operations[step_ops_start:], current_virtual_users, stability_time, now, @@ -288,7 +308,7 @@ async def run_load_step( step_start_time, stability_time, now, - operations, + operations[step_ops_start:], ): step_end_time = now break