diff --git a/src/sentry/issues/derived/tasks.py b/src/sentry/issues/derived/tasks.py index 5e4a4fbacebe..aa76c471cef0 100644 --- a/src/sentry/issues/derived/tasks.py +++ b/src/sentry/issues/derived/tasks.py @@ -581,14 +581,18 @@ def heal_stale_derived_data(**kwargs: object) -> None: lower_bound = 0 if stale_hash is None else state.stale[stale_hash] logger.info( "heal_stale_derived_data.range_selection_started", - extra={"hash_kind": hash_kind, "remaining_budget": remaining}, + extra={ + "hash_kind": hash_kind, + "remaining_budget": remaining, + "group_id_lower_bound": lower_bound, + }, ) range_selection_started_at = time.monotonic() try: range_result = group_id_ranges_for_hash( stale_hash, - chunk_size=batch_size, - max_chunks=remaining, + range_size=batch_size, + max_ranges=remaining, group_id_lower_bound=lower_bound, ) except OperationalError: @@ -597,6 +601,8 @@ def heal_stale_derived_data(**kwargs: object) -> None: extra={ "hash_kind": hash_kind, "elapsed": time.monotonic() - range_selection_started_at, + "pipeline_hash": stale_hash, + "group_id_lower_bound": lower_bound, }, ) metrics.incr( @@ -911,6 +917,7 @@ def regenerate_stale_derived_data_batch( ) from taskbroker_client.state import current_task + from sentry import options from sentry.issues.derived.promote import build_and_promote_batch from sentry.issues.models.groupderiveddata import GroupDerivedData from sentry.taskworker.selfchain_idempotency import already_spawned, mark_spawned @@ -944,6 +951,7 @@ def regenerate_stale_derived_data_batch( start = time.monotonic() + batch_size = max(1, options.get("issues.derived.heal-batch-size")) group_ids = list( GroupDerivedData.objects.filter( pipeline_hash=target_hash, # a None target renders as IS NULL @@ -951,7 +959,22 @@ def regenerate_stale_derived_data_batch( group_id__lt=group_id_end, ) .order_by("group_id") - .values_list("group_id", flat=True) + .values_list("group_id", flat=True)[: batch_size + 1] + ) + range_overflow = len(group_ids) > batch_size + if range_overflow: + group_ids = group_ids[:batch_size] + + # How close the scheduler's density estimate landed to reality. The worker + # records at most batch_size rows here and reports denser ranges separately as + # ``range_overflow`` reschedules. Counts well below batch_size mean ranges span + # too few IDs and scheduling slots are being wasted; frequent overflow means they + # span too many. Use both signals to tune the density sampling constants. + metrics.distribution( + "issues.derived.heal_range_rows_found", + len(group_ids), + sample_rate=1.0, + tags={"hash_kind": "null" if target_hash is None else "stale"}, ) result = build_and_promote_batch( @@ -981,6 +1004,21 @@ def regenerate_stale_derived_data_batch( ) if activation_id: mark_spawned(_REGENERATE_STALE_BATCH_TASK_KEY, activation_id) + elif range_overflow: + rescheduled = True + metrics.incr( + "issues.derived.regenerate_stale_batch_rescheduled", + sample_rate=1.0, + tags={"reason": "range_overflow"}, + ) + regenerate_stale_derived_data_batch.delay( + stale_pipeline_hashes=stale_pipeline_hashes, + target_hash=target_hash, + group_id_start=group_ids[-1] + 1, + group_id_end=group_id_end, + ) + if activation_id: + mark_spawned(_REGENERATE_STALE_BATCH_TASK_KEY, activation_id) _record_batch_metrics( result.processed, diff --git a/src/sentry/issues/derived/tasks_util.py b/src/sentry/issues/derived/tasks_util.py index 5da712e5c530..295c973c8099 100644 --- a/src/sentry/issues/derived/tasks_util.py +++ b/src/sentry/issues/derived/tasks_util.py @@ -2,12 +2,15 @@ import logging import random +import time +from collections.abc import Sequence from dataclasses import dataclass from datetime import datetime, timedelta, timezone +from math import ceil from typing import Protocol -from django.db import connections, router from django.db.models import Max, Min +from django.db.utils import OperationalError from sentry.issues.derived.check import CheckFailure, CheckId, CheckInvalidated, CheckResult from sentry.issues.models.groupderiveddata import GroupDerivedData @@ -18,11 +21,14 @@ logger = logging.getLogger(__name__) _MAX_CHECK_GROUPS = 10_000 - -# Safety valve on the number of group IDs one ``group_id_ranges_for_hash`` call may -# walk, however large the requested chunking is. -_MAX_SCANNED_GROUP_IDS = 2_000_000 -_GROUP_ID_RANGE_TIMEOUT = timedelta(seconds=50) +_GROUP_ID_RANGE_QUERY_TIMEOUT = timedelta(seconds=40) +# Use exact boundaries below this requested row count; estimate density above it. +_MAX_EXACT_RANGE_ROWS = 10_000 +# Each probe reads this many matching IDs and extrapolates the following number of ranges. +_RANGE_DENSITY_SAMPLE_SIZE = 100 +_RANGES_PER_DENSITY_SAMPLE = 5 +# Bound total probe work for regions with large scheduling budgets. +_MAX_RANGE_DENSITY_SAMPLES = 200 @dataclass(frozen=True) @@ -114,84 +120,147 @@ def _pick_random_fresh_group_ranges( return ranges +def _exact_group_id_ranges( + group_ids: Sequence[int], *, range_size: int, max_ranges: int +) -> list[tuple[int, int]]: + """Build exact ranges from ordered IDs, using an optional lookahead ID as the final end.""" + if not group_ids: + return [] + + requested_rows = range_size * max_ranges + starts = group_ids[:requested_rows:range_size] + last_end = group_ids[requested_rows] if len(group_ids) > requested_rows else group_ids[-1] + 1 + return list(zip(starts, [*starts[1:], last_end])) + + +def _estimate_group_id_ranges( + density_sample: Sequence[int], *, range_size: int, range_count: int +) -> list[tuple[int, int]]: + """Build contiguous ranges sized from the density of an ordered ID sample.""" + sample_width = density_sample[-1] - density_sample[0] + 1 + estimated_width = ceil(sample_width * range_size / len(density_sample)) + first_group_id = density_sample[0] + return [ + ( + first_group_id + index * estimated_width, + first_group_id + (index + 1) * estimated_width, + ) + for index in range(range_count) + ] + + def group_id_ranges_for_hash( - pipeline_hash: str | None, *, chunk_size: int, max_chunks: int, group_id_lower_bound: int = 0 + pipeline_hash: str | None, *, range_size: int, max_ranges: int, group_id_lower_bound: int = 0 ) -> GroupIdRangeResult: - """Partition the group IDs of GroupDerivedData rows with a pipeline_hash into ranges. + """Estimate ranges covering GroupDerivedData rows with a pipeline_hash. + + Returns at most max_ranges of ascending disjoint [start, end) ranges, each + targeting range_size group IDs. ``drained`` is true only when a valid query + found no rows at or above ``group_id_lower_bound``. + + Imagine a sequence of GroupDerivedData rows with some stale (s): - Returns at most max_chunks of ascending disjoint [start, end) ranges, each - covering chunk_size group IDs but possibly the last. ``drained`` is true only - when a valid query found no rows at or above ``group_id_lower_bound``. + s.......s.......s.......s...........ssssssss.... + 0 8 16 24 36..43 + + We can precisely query for ranges with an equal number of stale rows, but that requires + us to do a full index scan of those rows, and if we've been mutating, that can get + surprisingly slow, especially since we'd like to be able to divvy out 100s of thousands + of rows. Our range processing task also filters and is tolerant of variation, so instead + of trying to be exact, we approximate. + + Simply cutting the ID span into N equal ranges can be rough. Asking for 3 above gives + [0,16) [16,32) [32,48), holding 2, 2, and 8 stale rows: one range has two thirds of the + work, and that's bad for our goal of great throughput. + + Instead, we set a budget of how much we'd like to do, and take incremental 'core samples', + using each to size the ranges that follow it. Sampling 4 rows at a time and targeting 4 + rows per range, the first sample reads 0, 8, 16, 24, four rows spanning 25 IDs, so we emit + [0,25). The next sample resumes there, skips the empty gap entirely, and lands on + 36, 37, 38, 39, four rows spanning 4 IDs, so we emit [36,40). The last sample sees the end + of the data and falls back to exact boundaries, [40,44). + + That's 4, 4, 4 instead of 2, 2, 8, for 3 small queries instead of a full scan. It's still + an estimate: an over-dense range is split by the worker, an under-dense one costs only a + scheduling slot. In production a sample is _RANGE_DENSITY_SAMPLE_SIZE rows and sizes + _RANGES_PER_DENSITY_SAMPLE ranges, rather than one. + + Raises ``OperationalError`` if all density probes together exceed the query + budget. Callers must not interpret that as a drained hash. """ - if chunk_size <= 0 or max_chunks <= 0: + if range_size <= 0 or max_ranges <= 0: return GroupIdRangeResult(ranges=[], drained=False) - # One boundary per chunk, plus one extra to close the final range (or, if we ran - # out of matching rows first, to tell us we did). - scan_limit = chunk_size * (max_chunks + 1) - if scan_limit > _MAX_SCANNED_GROUP_IDS: - logger.warning( - "group_id_ranges_for_hash.scan_budget_clamped", - extra={ - "chunk_size": chunk_size, - "max_chunks": max_chunks, - "requested_scan_limit": scan_limit, - "max_scan_limit": _MAX_SCANNED_GROUP_IDS, - }, + matching_group_ids = ( + GroupDerivedData.objects.filter(pipeline_hash=pipeline_hash) + .order_by("pipeline_hash", "group_id") + .values_list("group_id", flat=True) + ) + using = matching_group_ids.db + matching_group_ids = matching_group_ids.using(using) + query_deadline = time.monotonic() + _GROUP_ID_RANGE_QUERY_TIMEOUT.total_seconds() + + with metrics.timer("issues.derived.group_id_range_query"): + + def fetch_group_ids(start: int, limit: int) -> list[int]: + remaining_seconds = query_deadline - time.monotonic() + if remaining_seconds <= 0.001: + raise OperationalError("group ID range query budget exceeded") + + with statement_timeout(using, timedelta(seconds=remaining_seconds)): + return list(matching_group_ids.filter(group_id__gte=start)[:limit]) + + requested_rows = range_size * max_ranges + if requested_rows <= _MAX_EXACT_RANGE_ROWS: + group_ids = fetch_group_ids(group_id_lower_bound, requested_rows + 1) + if not group_ids: + return GroupIdRangeResult(ranges=[], drained=True) + + return GroupIdRangeResult( + ranges=_exact_group_id_ranges( + group_ids, + range_size=range_size, + max_ranges=max_ranges, + ), + drained=False, + ) + + result_ranges: list[tuple[int, int]] = [] + next_group_id = group_id_lower_bound + density_sample_count = min( + ceil(max_ranges / _RANGES_PER_DENSITY_SAMPLE), + _MAX_RANGE_DENSITY_SAMPLES, ) - # Never clamp below two chunks: a lone boundary can't close a range, so we'd - # return nothing and the caller would take that to mean there was nothing to do. - scan_limit = max(_MAX_SCANNED_GROUP_IDS, chunk_size * 2) - - hash_predicate = "pipeline_hash IS NULL" if pipeline_hash is None else "pipeline_hash = %s" - # The LIMIT lives in the innermost subquery to encourage Postgres to enforce - # the limit before doing the window function stuff. ``cnt`` tells us whether we - # hit the limit, and pulls out the last row we saw so we can close the final - # chunk without a second query. - sql = f""" - SELECT group_id, rn, cnt FROM ( - SELECT group_id, - -- Ordered so the numbering is defined rather than dependent on the - -- order the subquery happens to emit. count(*) stays unordered; an - -- ORDER BY there would turn it into a running count. - row_number() OVER (ORDER BY group_id) - 1 AS rn, - count(*) OVER () AS cnt - FROM ( - SELECT group_id - FROM {GroupDerivedData._meta.db_table} - WHERE {hash_predicate} AND group_id >= %s - ORDER BY group_id - LIMIT %s - ) scanned - ) numbered - WHERE mod(rn, %s) = 0 OR rn = cnt - 1 - ORDER BY rn - """ - params: list[str | int] = [] if pipeline_hash is None else [pipeline_hash] - params += [group_id_lower_bound, scan_limit, chunk_size] - - using = router.db_for_read(GroupDerivedData) - with ( - statement_timeout(using, _GROUP_ID_RANGE_TIMEOUT), - metrics.timer("issues.derived.group_id_range_query"), - connections[using].cursor() as cursor, - ): - cursor.execute(sql, params) - rows = cursor.fetchall() - - if not rows: - return GroupIdRangeResult(ranges=[], drained=True) - - scanned = rows[0][2] - boundaries = [group_id for group_id, rn, _ in rows if rn % chunk_size == 0] - - if scanned == scan_limit: - # We got more than enough rows; the trailing boundary closes the last range. - ends = boundaries[1:] - else: - # We didn't fill out the final chunk, so the last row we scanned closes it. - ends = boundaries[1:] + [rows[-1][0] + 1] - return GroupIdRangeResult(ranges=list(zip(boundaries, ends))[:max_chunks], drained=False) + ranges_per_sample, samples_with_extra_range = divmod(max_ranges, density_sample_count) + # Ranges probably don't divide equally by samples, so we try to distribute the remainder + # cleanly. + range_counts = [ranges_per_sample + 1] * samples_with_extra_range + [ranges_per_sample] * ( + density_sample_count - samples_with_extra_range + ) + + for ranges_for_sample in range_counts: + sampled_group_ids = fetch_group_ids(next_group_id, _RANGE_DENSITY_SAMPLE_SIZE + 1) + if not sampled_group_ids: + return GroupIdRangeResult(ranges=result_ranges, drained=not result_ranges) + if len(sampled_group_ids) <= _RANGE_DENSITY_SAMPLE_SIZE: + # We requested one extra, so this means we've got all the data and can exit. + starts = sampled_group_ids[::range_size] + ends = starts[1:] + [sampled_group_ids[-1] + 1] + result_ranges.extend(zip(starts, ends)) + break + + # trim the over-query + density_sample = sampled_group_ids[:-1] + # generate range_for_sample ranges based on density_sample. + estimated_ranges = _estimate_group_id_ranges( + density_sample, + range_size=range_size, + range_count=ranges_for_sample, + ) + result_ranges.extend(estimated_ranges) + next_group_id = estimated_ranges[-1][1] + + return GroupIdRangeResult(ranges=result_ranges[:max_ranges], drained=False) def _resume_check_id( diff --git a/tests/sentry/issues/derived/test_tasks.py b/tests/sentry/issues/derived/test_tasks.py index 5769ce118f75..cd0cd82a4eca 100644 --- a/tests/sentry/issues/derived/test_tasks.py +++ b/tests/sentry/issues/derived/test_tasks.py @@ -31,6 +31,8 @@ from sentry.issues.derived.tasks_util import ( GroupIdRangeResult, SpawnState, + _estimate_group_id_ranges, + _exact_group_id_ranges, _pick_random_fresh_group_ranges, group_id_ranges_for_hash, ) @@ -645,6 +647,42 @@ def test_range_selection_timeout_is_reported_and_other_hashes_continue(self) -> in mock_incr.call_args_list ) + def test_range_selection_timeout_does_not_advance_mark(self) -> None: + stale = self._pick_stale_hash() + state = HealSchedulerState( + head_hash=PIPELINE.pipeline_hash, + stale={stale: 10}, + discovered_at=datetime.now(timezone.utc), + ) + with ( + override_options( + { + "issues.derived.heal-max-tasks": 1, + "issues.derived.check-task-count": 0, + } + ), + patch("sentry.issues.derived.tasks.load_state", return_value=state), + patch( + "sentry.issues.derived.tasks_util.group_id_ranges_for_hash", + side_effect=[ + GroupIdRangeResult(ranges=[], drained=True), + OperationalError, + ], + ), + patch.object(regenerate_stale_derived_data_batch, "delay") as delay, + patch("sentry.issues.derived.tasks.logger") as mock_logger, + ): + heal_stale_derived_data() + + assert state.stale == {stale: 10} + delay.assert_not_called() + failure_log = mock_logger.exception.call_args + assert failure_log.args == ("heal_stale_derived_data.range_selection_failed",) + assert failure_log.kwargs["extra"]["hash_kind"] == "stale" + assert failure_log.kwargs["extra"]["pipeline_hash"] == stale + assert failure_log.kwargs["extra"]["group_id_lower_bound"] == 10 + assert failure_log.kwargs["extra"]["elapsed"] >= 0 + def test_discovered_hashes_are_saved_before_range_selection(self) -> None: stale_hash = self._pick_stale_hash() with ( @@ -793,8 +831,8 @@ def test_null_starts_at_zero_without_rewriting_state(self) -> None: mock_ranges.assert_called_once_with( None, - chunk_size=500, - max_chunks=1, + range_size=500, + max_ranges=1, group_id_lower_bound=0, ) mock_save.assert_not_called() @@ -1374,6 +1412,26 @@ def test_caps_total_groups(self) -> None: assert result == [(group_ids[0], group_ids[1] + 1)] +class TestGroupIdRangeMath: + def test_exact_ranges_use_lookahead_to_close_final_range(self) -> None: + assert _exact_group_id_ranges([10, 20, 30, 40, 50], range_size=2, max_ranges=2) == [ + (10, 30), + (30, 50), + ] + + def test_exact_ranges_close_short_tail_after_last_group(self) -> None: + assert _exact_group_id_ranges([10, 20, 30], range_size=2, max_ranges=5) == [ + (10, 30), + (30, 31), + ] + + def test_estimated_ranges_use_sample_density(self) -> None: + assert _estimate_group_id_ranges([10, 20, 31], range_size=2, range_count=2) == [ + (10, 25), + (25, 40), + ] + + @with_feature("projects:issue-action-log-write-to-db") class GroupIdRangesForHashTest(DerivedDataTaskTestBase): HASH = "a" * 16 @@ -1389,55 +1447,54 @@ def test_no_matching_rows(self) -> None: self._seed(2, self.OTHER_HASH) assert group_id_ranges_for_hash( - self.HASH, chunk_size=2, max_chunks=5 + self.HASH, range_size=2, max_ranges=5 ) == GroupIdRangeResult(ranges=[], drained=True) - assert group_id_ranges_for_hash(None, chunk_size=2, max_chunks=5) == GroupIdRangeResult( + assert group_id_ranges_for_hash(None, range_size=2, max_ranges=5) == GroupIdRangeResult( ranges=[], drained=True ) def test_query_has_statement_timeout(self) -> None: with patch("sentry.issues.derived.tasks_util.statement_timeout") as timeout: - group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) + group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=5) - assert timeout.call_args.args[1] == timedelta(seconds=50) + query_timeout = timeout.call_args.args[1] + assert timedelta(0) < query_timeout <= timedelta(seconds=40) def test_short_tail_is_one_range(self) -> None: null_ids = self._seed(3, None) hash_ids = self._seed(3, self.HASH) - # Fewer rows than chunk_size, for both the NULL and the concrete-hash + # Fewer rows than range_size, for both the NULL and the concrete-hash # predicate, and neither picks up the other's rows. - assert group_id_ranges_for_hash(None, chunk_size=10, max_chunks=5).ranges == [ + assert group_id_ranges_for_hash(None, range_size=10, max_ranges=5).ranges == [ (null_ids[0], null_ids[-1] + 1) ] - assert group_id_ranges_for_hash(self.HASH, chunk_size=10, max_chunks=5).ranges == [ + assert group_id_ranges_for_hash(self.HASH, range_size=10, max_ranges=5).ranges == [ (hash_ids[0], hash_ids[-1] + 1) ] def test_exact_chunk_boundaries(self) -> None: group_ids = self._seed(5, self.HASH) - assert group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5).ranges == [ + assert group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=5).ranges == [ (group_ids[0], group_ids[2]), (group_ids[2], group_ids[4]), (group_ids[4], group_ids[4] + 1), ] - def test_truncates_to_max_chunks(self) -> None: + def test_truncates_to_max_ranges(self) -> None: group_ids = self._seed(5, self.HASH) # The 5th group is left out rather than folded into an oversized last range. - assert group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=2).ranges == [ + assert group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=2).ranges == [ (group_ids[0], group_ids[2]), (group_ids[2], group_ids[4]), ] - def test_truncates_when_scan_limit_is_reached(self) -> None: - # 6 rows exactly fills the scan budget of chunk_size * (max_chunks + 1), so - # the tail is known to be incomplete and waits for the next call. + def test_truncates_full_tail_to_max_ranges(self) -> None: group_ids = self._seed(6, self.HASH) - assert group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=2).ranges == [ + assert group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=2).ranges == [ (group_ids[0], group_ids[2]), (group_ids[2], group_ids[4]), ] @@ -1446,33 +1503,61 @@ def test_invalid_chunking(self) -> None: self._seed(2, self.HASH) assert group_id_ranges_for_hash( - self.HASH, chunk_size=0, max_chunks=5 + self.HASH, range_size=0, max_ranges=5 ) == GroupIdRangeResult(ranges=[], drained=False) assert group_id_ranges_for_hash( - self.HASH, chunk_size=2, max_chunks=0 + self.HASH, range_size=2, max_ranges=0 ) == GroupIdRangeResult(ranges=[], drained=False) - def test_clamps_scan_budget(self) -> None: - group_ids = self._seed(3, self.HASH) + def test_density_probes_share_one_query_budget(self) -> None: + with ( + patch( + "sentry.issues.derived.tasks_util.time.monotonic", + # Query deadline, metrics timer start, budget check, metrics timer end. + side_effect=[0.0, 0.0, 41.0, 41.0], + ), + pytest.raises(OperationalError, match="query budget exceeded"), + ): + group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=5) + + def test_samples_local_density_for_approximate_ranges(self) -> None: + groups = self.create_unprocessed_groups(81) + all_group_ids = sorted(group.id for group in groups) + matching_group_ids = [all_group_ids[index] for index in range(0, 81, 10)] + for group_id in all_group_ids: + GroupDerivedData.objects.create( + group_id=group_id, + pipeline_hash=self.HASH if group_id in matching_group_ids else self.OTHER_HASH, + ) - with patch("sentry.issues.derived.tasks_util._MAX_SCANNED_GROUP_IDS", 2): - result = group_id_ranges_for_hash(self.HASH, chunk_size=1, max_chunks=5) + with ( + patch("sentry.issues.derived.tasks_util._MAX_EXACT_RANGE_ROWS", 0), + patch("sentry.issues.derived.tasks_util._RANGE_DENSITY_SAMPLE_SIZE", 2), + patch("sentry.issues.derived.tasks_util._RANGES_PER_DENSITY_SAMPLE", 2), + patch("sentry.issues.derived.tasks_util._MAX_RANGE_DENSITY_SAMPLES", 2), + ): + result = group_id_ranges_for_hash(self.HASH, range_size=4, max_ranges=4) - # Only 2 rows were scanned, so the clamp must not be mistaken for having - # reached the end — the tail may not run past what we scanned. - assert result.ranges == [(group_ids[0], group_ids[1])] + assert len(result.ranges) == 4 + assert result.ranges == sorted(result.ranges) + assert all( + any(start <= group_id < end for start, end in result.ranges) + for group_id in matching_group_ids + ) - def test_clamp_never_drops_below_one_chunk(self) -> None: - group_ids = self._seed(3, self.HASH) + def test_low_volume_request_uses_exact_boundaries(self) -> None: + group_ids = self._seed(5, self.HASH) - # Clamping below chunk_size would leave a single boundary, which can't close - # a range — the caller would read the empty result as "nothing to do". - with patch("sentry.issues.derived.tasks_util._MAX_SCANNED_GROUP_IDS", 1): - result = group_id_ranges_for_hash(self.HASH, chunk_size=2, max_chunks=5) + with ( + patch("sentry.issues.derived.tasks_util._MAX_EXACT_RANGE_ROWS", 10), + patch("sentry.issues.derived.tasks_util._RANGE_DENSITY_SAMPLE_SIZE", 2), + ): + result = group_id_ranges_for_hash(self.HASH, range_size=2, max_ranges=5) assert result.ranges == [ (group_ids[0], group_ids[2]), - (group_ids[2], group_ids[2] + 1), + (group_ids[2], group_ids[4]), + (group_ids[4], group_ids[4] + 1), ] def test_lower_bound(self) -> None: @@ -1480,8 +1565,8 @@ def test_lower_bound(self) -> None: result = group_id_ranges_for_hash( self.HASH, - chunk_size=2, - max_chunks=5, + range_size=2, + max_ranges=5, group_id_lower_bound=group_ids[1], ) @@ -1642,6 +1727,42 @@ def test_reschedules_on_batch_timeout(self) -> None: assert kwargs["group_id_start"] == group_ids[0] + 1 assert kwargs["group_id_end"] == group_ids[-1] + 1 + def test_reschedules_when_estimated_range_is_overfull(self) -> None: + groups = self.create_unprocessed_groups(3) + group_ids = sorted(g.id for g in groups) + for gid in group_ids: + process_group_log(gid) + + stale = self._stale() + GroupDerivedData.objects.filter(group_id__in=group_ids).update(pipeline_hash=stale) + + with ( + override_options({"issues.derived.heal-batch-size": 2}), + patch.object(regenerate_stale_derived_data_batch, "delay") as mock_delay, + ): + regenerate_stale_derived_data_batch( + stale_pipeline_hashes=[stale], + target_hash=stale, + group_id_start=group_ids[0], + group_id_end=group_ids[-1] + 1, + ) + + assert ( + GroupDerivedData.objects.get(group_id=group_ids[0]).pipeline_hash + == PIPELINE.pipeline_hash + ) + assert ( + GroupDerivedData.objects.get(group_id=group_ids[1]).pipeline_hash + == PIPELINE.pipeline_hash + ) + assert GroupDerivedData.objects.get(group_id=group_ids[2]).pipeline_hash == stale + mock_delay.assert_called_once_with( + stale_pipeline_hashes=[stale], + target_hash=stale, + group_id_start=group_ids[1] + 1, + group_id_end=group_ids[-1] + 1, + ) + def test_reschedules_on_group_log_timeout(self) -> None: groups = self.create_unprocessed_groups(2) group_ids = sorted(g.id for g in groups)