Skip to content
Closed
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 25 additions & 5 deletions monitoring/benchmarker/engine/loads/step_execution.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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."""

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I don't understand why the operations are ordered within a 60 seconds tolerance because they are recorded as they complete?
Also: I think making this assumption here of logic defined elsewhere is not good. If that assumption changes (and I don't think it is guaranteed by anything?), this will break this logic here.
Maybe replace operations with a data structure that has some guarantees on the sorting?



@dataclass
class ActiveUser:
Expand Down Expand Up @@ -198,6 +201,21 @@ async def wind_down_and_cleanup_remaining_users(
)


def first_index_completed_since(

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Have you considered instead doing something like summarize_and_report_step?

    step_ops = [
        op
        for op in operations
        if op.completed_at.datetime >= step_start_time
        and op.completed_at.datetime <= step_end_time
    ]

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,
Expand All @@ -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
Expand All @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand Down
Loading