From 19a5e183d84b10b7c8912c226465addf73e274f5 Mon Sep 17 00:00:00 2001 From: Kyle Consalus Date: Thu, 6 Aug 2026 22:13:19 -0700 Subject: [PATCH 1/2] Checker --- src/sentry/issues/derived/check.py | 155 +++++++++++ src/sentry/issues/derived/tasks.py | 222 +++++++++++++++ src/sentry/issues/derived/tasks_util.py | 82 ++++++ src/sentry/options/defaults.py | 7 + .../sentry/issues/derived/test_processing.py | 129 +++++++++ tests/sentry/issues/derived/test_tasks.py | 262 +++++++++++++++++- 6 files changed, 854 insertions(+), 3 deletions(-) create mode 100644 src/sentry/issues/derived/check.py create mode 100644 src/sentry/issues/derived/tasks_util.py diff --git a/src/sentry/issues/derived/check.py b/src/sentry/issues/derived/check.py new file mode 100644 index 000000000000..d5818a8f3595 --- /dev/null +++ b/src/sentry/issues/derived/check.py @@ -0,0 +1,155 @@ +import time +from dataclasses import dataclass +from datetime import datetime, timedelta +from enum import Enum +from typing import NamedTuple +from uuid import uuid4 + +from sentry.issues.derived.processing import DEFAULT_BATCH_SIZE, PIPELINE +from sentry.issues.derived.store import GroupDerivedDataStore +from sentry.issues.models.groupactionlogentry import GroupActionLogEntry +from sentry.issues.models.groupderiveddata import EPOCH, GroupDerivedData +from sentry.workflow_engine.caches.mapping import CacheMapping + + +class CheckSuccess(Enum): + OK = "ok" + + +@dataclass(frozen=True) +class CheckFailure: + group_id: int + cursor_date: datetime + cursor_id: int + features: frozenset[str] + + +type CheckResult = CheckSuccess | CheckFailure + + +class CheckId(NamedTuple): + invocation_id: str + group_id: int + generated_at: datetime + cursor_date: datetime + cursor_id: int + pipeline_hash: str + + +_check_cache = CacheMapping[CheckId, GroupDerivedData]( + lambda key: ( + f"{key.invocation_id}:{key.group_id}:{key.generated_at.isoformat()}:" + f"{key.cursor_date.isoformat()}:{key.cursor_id}:{key.pipeline_hash}" + ), + namespace="gdd-check", + ttl_seconds=86400, +) + + +class CheckTimeout(Exception): + def __init__(self, check_id: CheckId) -> None: + self.check_id = check_id + super().__init__(check_id) + + +def _entries_after_cursor( + derived: GroupDerivedData, target: CheckId, batch_size: int +) -> list[GroupActionLogEntry]: + return list( + GroupActionLogEntry.objects.filter(group_id=target.group_id) + .extra( + where=[ + 'ROW("date_added", "id") > ROW(%s, %s)', + 'ROW("date_added", "id") <= ROW(%s, %s)', + ], + params=[ + derived.cursor_date, + derived.cursor_id, + target.cursor_date, + target.cursor_id, + ], + ) + .order_by("date_added", "id")[:batch_size] + ) + + +def check_derived_data( + derived: GroupDerivedData, + timeout: timedelta | None = None, + *, + check_id: CheckId | None = None, + batch_size: int = DEFAULT_BATCH_SIZE, +) -> CheckResult | None: + if derived.pipeline_hash != PIPELINE.pipeline_hash: + return None + + target_fields = ( + derived.group_id, + derived.generated_at, + derived.cursor_date, + derived.cursor_id, + derived.pipeline_hash, + ) + if check_id is not None: + if check_id[1:] != target_fields: + _check_cache.delete(check_id) + return None + target = check_id + else: + target = CheckId(uuid4().hex, *target_fields) + + replayed_derived = _check_cache.get(target) if check_id is not None else None + if replayed_derived is None: + replayed_derived = GroupDerivedData( + group_id=derived.group_id, + cursor_date=EPOCH, + cursor_id=0, + data={}, + pipeline_hash=derived.pipeline_hash, + ) + + deadline = time.monotonic() + timeout.total_seconds() if timeout is not None else None + while entries := _entries_after_cursor(replayed_derived, target, batch_size): + replayed = PIPELINE.run( + entries, state=GroupDerivedDataStore.load(PIPELINE, replayed_derived) + ) + GroupDerivedDataStore.apply_to_instance( + replayed_derived, GroupDerivedDataStore.build_update(PIPELINE, replayed) + ) + replayed_derived.cursor_date = entries[-1].date_added + replayed_derived.cursor_id = entries[-1].id + if len(entries) < batch_size: + break + if deadline is not None and time.monotonic() >= deadline: + _check_cache.set(target, replayed_derived) + raise CheckTimeout(target) + + current = ( + GroupDerivedData.objects.filter(group_id=target.group_id) + .values_list("generated_at", "cursor_date", "cursor_id", "pipeline_hash") + .first() + ) + if current != ( + target.generated_at, + target.cursor_date, + target.cursor_id, + target.pipeline_hash, + ): + _check_cache.delete(target) + return None + + replayed = GroupDerivedDataStore.load(PIPELINE, replayed_derived) + stored = GroupDerivedDataStore.load(PIPELINE, derived) + different_features = frozenset( + feature.name for feature in PIPELINE.features if replayed[feature] != stored[feature] + ) + _check_cache.delete(target) + if not different_features: + return CheckSuccess.OK + + return CheckFailure( + group_id=derived.group_id, + cursor_date=derived.cursor_date, + cursor_id=derived.cursor_id, + features=different_features, + ) diff --git a/src/sentry/issues/derived/tasks.py b/src/sentry/issues/derived/tasks.py index d36317a9ac4e..c2f1138d7007 100644 --- a/src/sentry/issues/derived/tasks.py +++ b/src/sentry/issues/derived/tasks.py @@ -27,10 +27,13 @@ _GENERATE_PROJECT_TASK_KEY = "generate_project_derived_data" _GENERATE_BATCH_TASK_KEY = "generate_project_derived_data_batch" _REGENERATE_STALE_BATCH_TASK_KEY = "regenerate_stale_derived_data_batch" +_CHECK_FRESH_BATCH_TASK_KEY = "check_fresh_derived_data_batch" _GENERATE_GROUP_TASK_KEY = "generate_group_derived_data" +_CHECK_GROUP_TASK_KEY = "check_group_derived_data" # Cap self-rescheduling rebuilds to avoid infinite loops on very large groups. _MAX_GENERATION_RUNS = 20 +_MAX_CHECK_RUNS = 20 # Maximum group IDs loaded by one project-level task invocation. _MAX_PROJECT_GROUPS = 10_000 # Hard cap on distinct stale pipeline hashes handled per heal invocation. @@ -207,6 +210,93 @@ def process_group_log_task(group_id: int, incremental: bool = False, **kwargs: o logger.info("process_group_log_task.group_not_found", extra={"group_id": group_id}) +@instrumented_task( + name="sentry.issues.derived.tasks.check_group_derived_data", + namespace=issues_tasks, + silo_mode=SiloMode.CELL, + processing_deadline_duration=int(BATCH_PROCESSING_DEADLINE.total_seconds()), +) +def check_group_derived_data( + group_id: int, + resume_check_id: str | None = None, + resume_generated_at: str | None = None, + resume_cursor_date: str | None = None, + resume_cursor_id: int | None = None, + resume_pipeline_hash: str | None = None, + prior_runs: int = 0, + **kwargs: object, +) -> None: + from taskbroker_client.state import current_task + + from sentry.issues.derived.check import CheckTimeout, check_derived_data + from sentry.issues.derived.tasks_util import _record_check_result, _resume_check_id + from sentry.issues.models.groupderiveddata import GroupDerivedData + from sentry.taskworker.selfchain_idempotency import already_spawned, mark_spawned + + task_state = current_task() + activation_id = task_state.id if task_state else None + if activation_id and already_spawned(_CHECK_GROUP_TASK_KEY, activation_id): + logger.info( + "check_group_derived_data.duplicate_skipped", + extra={"group_id": group_id, "activation_id": activation_id}, + ) + metrics.incr( + "taskworker.selfchain.duplicate_skipped", + tags={"task": _CHECK_GROUP_TASK_KEY}, + ) + return + + check_id = _resume_check_id( + group_id, + resume_check_id, + resume_generated_at, + resume_cursor_date, + resume_cursor_id, + resume_pipeline_hash, + ) + + derived = GroupDerivedData.objects.filter(group_id=group_id).first() + if derived is None: + _record_check_result(None) + return + + try: + result = check_derived_data( + derived, + timeout=BATCH_RETRIGGER_TIMEOUT, + check_id=check_id, + ) + except CheckTimeout as error: + if prior_runs + 1 >= _MAX_CHECK_RUNS: + logger.error( + "check_group_derived_data.max_runs_exceeded", + extra={ + "group_id": group_id, + "check_id": error.check_id, + "prior_runs": prior_runs + 1, + }, + ) + metrics.incr( + "issues.derived.check_group", sample_rate=1.0, tags={"result": "no_result"} + ) + return + + check_group_derived_data.delay( + group_id, + resume_check_id=error.check_id.invocation_id, + resume_generated_at=error.check_id.generated_at.isoformat(), + resume_cursor_date=error.check_id.cursor_date.isoformat(), + resume_cursor_id=error.check_id.cursor_id, + resume_pipeline_hash=error.check_id.pipeline_hash, + prior_runs=prior_runs + 1, + ) + if activation_id: + mark_spawned(_CHECK_GROUP_TASK_KEY, activation_id) + return + + _record_check_result(result) + + @instrumented_task( name="sentry.issues.derived.tasks.generate_group_derived_data", namespace=issues_tasks, @@ -531,6 +621,7 @@ def heal_stale_derived_data(**kwargs: object) -> None: """Rebuild a chunk of GroupDerivedData rows whose ``pipeline_hash`` is stale/NULL.""" from sentry import options from sentry.issues.derived.processing import PIPELINE + from sentry.issues.derived.tasks_util import _pick_random_fresh_group_range from sentry.issues.models.groupderiveddata import GroupDerivedData if not options.get("issues.derived.heal-enabled"): @@ -553,6 +644,23 @@ def heal_stale_derived_data(**kwargs: object) -> None: ) if not group_ids: logger.info("heal_stale_derived_data.nothing_to_heal") + check_range = _pick_random_fresh_group_range( + current_hash, options.get("issues.derived.check-sample-size") + ) + if check_range is not None: + start, end = check_range + check_fresh_derived_data_batch.delay( + group_id_start=start, + group_id_end=end, + ) + + logger.info( + "heal_stale_derived_data.checks_scheduled", + extra={ + "scheduled": check_range is not None, + "pipeline_hash": current_hash, + }, + ) return ranges = _chunk_group_ids_into_ranges(group_ids, batch_size)[:max_tasks] @@ -575,6 +683,120 @@ def heal_stale_derived_data(**kwargs: object) -> None: ) +@instrumented_task( + name="sentry.issues.derived.tasks.check_fresh_derived_data_batch", + namespace=issues_tasks, + silo_mode=SiloMode.CELL, + processing_deadline_duration=int(BATCH_PROCESSING_DEADLINE.total_seconds()), +) +def check_fresh_derived_data_batch( + group_id_start: int, + group_id_end: int, + resume_check_id: str | None = None, + resume_generated_at: str | None = None, + resume_cursor_date: str | None = None, + resume_cursor_id: int | None = None, + resume_pipeline_hash: str | None = None, + prior_runs: int = 0, + **kwargs: object, +) -> None: + """Check fresh GroupDerivedData rows in ``[group_id_start, group_id_end)``.""" + from taskbroker_client.state import current_task + + from sentry.issues.derived.check import CheckTimeout, check_derived_data + from sentry.issues.derived.processing import PIPELINE + from sentry.issues.derived.tasks_util import _record_check_result, _resume_check_id + from sentry.issues.models.groupderiveddata import GroupDerivedData + from sentry.taskworker.selfchain_idempotency import already_spawned, mark_spawned + + task_state = current_task() + activation_id = task_state.id if task_state else None + if activation_id and already_spawned(_CHECK_FRESH_BATCH_TASK_KEY, activation_id): + logger.info( + "check_fresh_derived_data_batch.duplicate_skipped", + extra={"group_id_start": group_id_start, "activation_id": activation_id}, + ) + metrics.incr( + "taskworker.selfchain.duplicate_skipped", + tags={"task": _CHECK_FRESH_BATCH_TASK_KEY}, + ) + return + + check_id = _resume_check_id( + group_id_start, + resume_check_id, + resume_generated_at, + resume_cursor_date, + resume_cursor_id, + resume_pipeline_hash, + ) + + derived_rows = GroupDerivedData.objects.filter( + pipeline_hash=PIPELINE.pipeline_hash, + group_id__gte=group_id_start, + group_id__lt=group_id_end, + ).order_by("group_id") + start = time.monotonic() + timeout_seconds = BATCH_RETRIGGER_TIMEOUT.total_seconds() + for derived in derived_rows.iterator(): + remaining = timedelta(seconds=max(0, timeout_seconds - (time.monotonic() - start))) + try: + result = check_derived_data( + derived, + timeout=remaining, + check_id=(check_id if derived.group_id == group_id_start else None), + ) + except CheckTimeout as error: + group_prior_runs = prior_runs if derived.group_id == group_id_start else 0 + if group_prior_runs + 1 >= _MAX_CHECK_RUNS: + logger.error( + "check_fresh_derived_data_batch.max_runs_exceeded", + extra={"group_id": derived.group_id, "check_id": error.check_id}, + ) + _record_check_result(None) + check_fresh_derived_data_batch.delay( + group_id_start=derived.group_id + 1, + group_id_end=group_id_end, + ) + if activation_id: + mark_spawned(_CHECK_FRESH_BATCH_TASK_KEY, activation_id) + return + + check_fresh_derived_data_batch.delay( + group_id_start=derived.group_id, + group_id_end=group_id_end, + resume_check_id=error.check_id.invocation_id, + resume_generated_at=error.check_id.generated_at.isoformat(), + resume_cursor_date=error.check_id.cursor_date.isoformat(), + resume_cursor_id=error.check_id.cursor_id, + resume_pipeline_hash=error.check_id.pipeline_hash, + prior_runs=group_prior_runs + 1, + ) + metrics.incr( + "issues.derived.check_fresh_batch_rescheduled", + sample_rate=1.0, + tags={"reason": "group_timeout"}, + ) + if activation_id: + mark_spawned(_CHECK_FRESH_BATCH_TASK_KEY, activation_id) + return + + _record_check_result(result) + if time.monotonic() - start >= timeout_seconds: + check_fresh_derived_data_batch.delay( + group_id_start=derived.group_id + 1, + group_id_end=group_id_end, + ) + metrics.incr( + "issues.derived.check_fresh_batch_rescheduled", + sample_rate=1.0, + tags={"reason": "batch_timeout"}, + ) + if activation_id: + mark_spawned(_CHECK_FRESH_BATCH_TASK_KEY, activation_id) + return + + @instrumented_task( name="sentry.issues.derived.tasks.regenerate_stale_derived_data_batch", namespace=issues_tasks, diff --git a/src/sentry/issues/derived/tasks_util.py b/src/sentry/issues/derived/tasks_util.py new file mode 100644 index 000000000000..3be3585bd9a1 --- /dev/null +++ b/src/sentry/issues/derived/tasks_util.py @@ -0,0 +1,82 @@ +import logging +import random +from datetime import datetime, timezone + +from django.db.models import Max, Min + +from sentry.issues.derived.check import CheckFailure, CheckId, CheckResult +from sentry.issues.models.groupderiveddata import GroupDerivedData +from sentry.utils import metrics + +logger = logging.getLogger(__name__) + +_MAX_CHECK_GROUPS = 10_000 + + +def _record_check_result(result: CheckResult | None) -> None: + outcome = "no_result" if result is None else "success" + if isinstance(result, CheckFailure): + outcome = "mismatch" + logger.warning( + "check_group_derived_data.mismatch", + extra={ + "group_id": result.group_id, + "cursor_date": result.cursor_date.isoformat(), + "cursor_id": result.cursor_id, + "features": sorted(result.features), + }, + ) + metrics.incr( + "issues.derived.check_group", + sample_rate=1.0, + tags={"result": outcome}, + ) + + +def _pick_random_fresh_group_range(pipeline_hash: str, target_size: int) -> tuple[int, int] | None: + """Pick a random range containing up to ``target_size`` fresh rows.""" + target_size = min(target_size, _MAX_CHECK_GROUPS) + if target_size <= 0: + return None + + fresh = GroupDerivedData.objects.filter(pipeline_hash=pipeline_hash) + bounds = fresh.aggregate(min_group_id=Min("group_id"), max_group_id=Max("group_id")) + min_group_id = bounds["min_group_id"] + max_group_id = bounds["max_group_id"] + if min_group_id is None or max_group_id is None: + return None + + random_start = random.randint(min_group_id, max_group_id) + group_ids = list( + fresh.filter(group_id__gte=random_start) + .order_by("group_id") + .values_list("group_id", flat=True)[:target_size] + ) + if not group_ids: + return None + return group_ids[0], group_ids[-1] + 1 + + +def _resume_check_id( + group_id: int, + invocation_id: str | None, + generated_at: str | None, + cursor_date: str | None, + cursor_id: int | None, + pipeline_hash: str | None, +) -> CheckId | None: + if None in (invocation_id, generated_at, cursor_date, cursor_id, pipeline_hash): + return None + assert invocation_id is not None + assert generated_at is not None + assert cursor_date is not None + assert cursor_id is not None + assert pipeline_hash is not None + return CheckId( + invocation_id, + group_id, + datetime.fromisoformat(generated_at).replace(tzinfo=timezone.utc), + datetime.fromisoformat(cursor_date).replace(tzinfo=timezone.utc), + cursor_id, + pipeline_hash, + ) diff --git a/src/sentry/options/defaults.py b/src/sentry/options/defaults.py index dbd4455b6cb1..b5a93484ad63 100644 --- a/src/sentry/options/defaults.py +++ b/src/sentry/options/defaults.py @@ -4113,6 +4113,13 @@ type=Int, flags=FLAG_AUTOMATOR_MODIFIABLE, ) +# Approximate number of fresh groups to check when there is no stale derived data to heal. +register( + "issues.derived.check-sample-size", + default=1000, + type=Int, + flags=FLAG_AUTOMATOR_MODIFIABLE, +) # Kill switch for Objectstore Debug Files migration register( diff --git a/tests/sentry/issues/derived/test_processing.py b/tests/sentry/issues/derived/test_processing.py index fe0531e96f88..5f5ac80784b7 100644 --- a/tests/sentry/issues/derived/test_processing.py +++ b/tests/sentry/issues/derived/test_processing.py @@ -1,3 +1,4 @@ +from collections.abc import Iterable from datetime import datetime, timedelta, timezone from unittest.mock import patch @@ -25,6 +26,7 @@ ) from sentry.issues.derived import processing from sentry.issues.derived.aggregators import AGGREGATORS +from sentry.issues.derived.check import CheckFailure, CheckSuccess, CheckTimeout, check_derived_data from sentry.issues.derived.features import ( BLOCKER, HAS_OPEN_FIX_PR, @@ -39,6 +41,7 @@ AggregatorResult, Feature, Pipeline, + State, StateUpdate, StateView, aggregator, @@ -131,6 +134,132 @@ def test_noop_when_no_new_entries(self) -> None: derived = process_group_log(group.id) assert derived.date_updated == old_updated + def test_check_derived_data_matches_replayed_state(self) -> None: + group = self.create_group() + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + derived = process_group_log(group.id) + + assert check_derived_data(derived) is CheckSuccess.OK + + def test_check_derived_data_reports_different_features(self) -> None: + group = self.create_group() + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + derived = process_group_log(group.id) + derived.view_count = 0 + + assert check_derived_data(derived) == CheckFailure( + group_id=group.id, + cursor_date=derived.cursor_date, + cursor_id=derived.cursor_id, + features=frozenset({VIEW_COUNT.name}), + ) + + def test_check_derived_data_skips_stale_pipeline(self) -> None: + group = self.create_group() + derived = process_group_log(group.id) + derived.pipeline_hash = "stale" + + assert check_derived_data(derived) is None + + def test_check_derived_data_can_resume_after_timeout(self) -> None: + group = self.create_group() + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + derived = process_group_log(group.id) + + with pytest.raises(CheckTimeout) as exc_info: + check_derived_data(derived, timeout=timedelta(0), batch_size=1) + + assert ( + check_derived_data( + derived, + timeout=timedelta(minutes=1), + check_id=exc_info.value.check_id, + batch_size=1, + ) + is CheckSuccess.OK + ) + + def test_check_derived_data_uses_invocation_scoped_checkpoints(self) -> None: + group = self.create_group() + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + derived = process_group_log(group.id) + + with pytest.raises(CheckTimeout) as first_timeout: + check_derived_data(derived, timeout=timedelta(0), batch_size=1) + with pytest.raises(CheckTimeout) as second_timeout: + check_derived_data(derived, timeout=timedelta(0), batch_size=1) + + assert ( + first_timeout.value.check_id.invocation_id + != second_timeout.value.check_id.invocation_id + ) + + for check_id in (first_timeout.value.check_id, second_timeout.value.check_id): + assert ( + check_derived_data( + derived, + timeout=timedelta(minutes=1), + check_id=check_id, + batch_size=1, + ) + is CheckSuccess.OK + ) + + def test_check_derived_data_stops_after_partial_batch(self) -> None: + group = self.create_group() + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + derived = process_group_log(group.id) + entries = list(GroupActionLogEntry.objects.filter(group_id=group.id)) + + with patch( + "sentry.issues.derived.check._entries_after_cursor", side_effect=[entries] + ) as get: + assert check_derived_data(derived, batch_size=2) is CheckSuccess.OK + + get.assert_called_once() + + def test_check_derived_data_does_not_resume_across_generations(self) -> None: + group = self.create_group() + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + derived = process_group_log(group.id) + + with pytest.raises(CheckTimeout) as exc_info: + check_derived_data(derived, timeout=timedelta(0), batch_size=1) + + GroupDerivedData.objects.filter(group_id=group.id).update( + generated_at=derived.generated_at + timedelta(seconds=1) + ) + derived.refresh_from_db() + + assert ( + check_derived_data( + derived, + timeout=timedelta(minutes=1), + check_id=exc_info.value.check_id, + batch_size=1, + ) + is None + ) + + def test_check_derived_data_skips_row_invalidated_during_replay(self) -> None: + group = self.create_group() + _publish(group=group, action=ViewAction(), actor=GroupActionActor.user(self.user.id)) + derived = process_group_log(group.id) + original_run = PIPELINE.run + + def run_and_invalidate( + entries: Iterable[GroupActionLogEntry], state: State | None = None + ) -> State: + result = original_run(entries, state=state) + GroupDerivedData.objects.filter(group_id=group.id).update(pipeline_hash=None) + return result + + with patch.object(PIPELINE, "run", side_effect=run_and_invalidate): + assert check_derived_data(derived) is None + def test_process_group_log_only_affects_target(self) -> None: group_a = self.create_group() group_b = self.create_group() diff --git a/tests/sentry/issues/derived/test_tasks.py b/tests/sentry/issues/derived/test_tasks.py index 884621ad329b..571ad8da4b74 100644 --- a/tests/sentry/issues/derived/test_tasks.py +++ b/tests/sentry/issues/derived/test_tasks.py @@ -1,18 +1,22 @@ from collections.abc import Sequence from datetime import datetime, timezone -from unittest.mock import patch +from unittest.mock import call, patch from sentry.issues.action_log.publish import publish_action from sentry.issues.action_log.types import ActionSource, GroupActionActor, ViewAction +from sentry.issues.derived.check import CheckId, CheckTimeout from sentry.issues.derived.processing import PIPELINE, GroupLogTimeout, process_group_log from sentry.issues.derived.tasks import ( BATCH_RETRIGGER_TIMEOUT, _discover_stale_pipeline_hashes, + check_fresh_derived_data_batch, + check_group_derived_data, generate_project_derived_data, generate_project_derived_data_batch, heal_stale_derived_data, regenerate_stale_derived_data_batch, ) +from sentry.issues.derived.tasks_util import _pick_random_fresh_group_range from sentry.issues.models.groupderiveddata import GroupDerivedData from sentry.models.group import Group from sentry.testutils.cases import TestCase @@ -40,6 +44,95 @@ def create_unprocessed_groups(self, count: int) -> list[Group]: return groups +@with_feature("projects:issue-action-log-write-to-db") +class CheckGroupDerivedDataTest(DerivedDataTaskTestBase): + def test_records_no_result_without_derived_data(self) -> None: + group = self.create_group(project=self.project) + + with patch("sentry.issues.derived.tasks_util.metrics.incr") as mock_incr: + check_group_derived_data(group.id) + + mock_incr.assert_called_once_with( + "issues.derived.check_group", + sample_rate=1.0, + tags={"result": "no_result"}, + ) + + def test_records_success(self) -> None: + group = self.create_unprocessed_groups(1)[0] + process_group_log(group.id) + + with patch("sentry.issues.derived.tasks_util.metrics.incr") as mock_incr: + check_group_derived_data(group.id) + + mock_incr.assert_called_once_with( + "issues.derived.check_group", + sample_rate=1.0, + tags={"result": "success"}, + ) + + def test_records_and_logs_mismatch(self) -> None: + group = self.create_unprocessed_groups(1)[0] + process_group_log(group.id) + GroupDerivedData.objects.filter(group_id=group.id).update(view_count=0) + + with ( + patch("sentry.issues.derived.tasks_util.metrics.incr") as mock_incr, + patch("sentry.issues.derived.tasks_util.logger.warning") as mock_warning, + ): + check_group_derived_data(group.id) + + derived = GroupDerivedData.objects.get(group_id=group.id) + mock_incr.assert_called_once_with( + "issues.derived.check_group", + sample_rate=1.0, + tags={"result": "mismatch"}, + ) + assert mock_warning.call_args == call( + "check_group_derived_data.mismatch", + extra={ + "group_id": group.id, + "cursor_date": derived.cursor_date.isoformat(), + "cursor_id": derived.cursor_id, + "features": ["view_count"], + }, + ) + + def test_reschedules_timeout_with_check_id(self) -> None: + group = self.create_unprocessed_groups(1)[0] + derived = process_group_log(group.id) + assert derived.pipeline_hash is not None + check_id = CheckId( + "invocation-id", + group.id, + derived.generated_at, + derived.cursor_date, + derived.cursor_id, + derived.pipeline_hash, + ) + + with ( + patch( + "sentry.issues.derived.check.check_derived_data", + side_effect=CheckTimeout(check_id), + ), + patch.object(check_group_derived_data, "delay") as mock_delay, + patch("sentry.issues.derived.tasks.metrics.incr") as mock_incr, + ): + check_group_derived_data(group.id) + + mock_delay.assert_called_once_with( + group.id, + resume_check_id="invocation-id", + resume_generated_at=derived.generated_at.isoformat(), + resume_cursor_date=derived.cursor_date.isoformat(), + resume_cursor_id=derived.cursor_id, + resume_pipeline_hash=derived.pipeline_hash, + prior_runs=1, + ) + mock_incr.assert_not_called() + + @with_feature("projects:issue-action-log-write-to-db") class GenerateProjectDerivedDataStaleOnlyTest(DerivedDataTaskTestBase): def test_only_includes_stale_groups(self) -> None: @@ -259,10 +352,37 @@ def test_no_stale_data(self) -> None: for g in groups: process_group_log(g.id) - with patch.object(regenerate_stale_derived_data_batch, "delay") as mock_delay: + group_ids = sorted(group.id for group in groups) + with ( + patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[0]), + patch.object(regenerate_stale_derived_data_batch, "delay") as mock_regenerate, + patch.object(check_fresh_derived_data_batch, "delay") as mock_check, + ): heal_stale_derived_data() - mock_delay.assert_not_called() + mock_regenerate.assert_not_called() + mock_check.assert_called_once_with( + group_id_start=group_ids[0], + group_id_end=group_ids[-1] + 1, + ) + + def test_checks_random_configured_group_range(self) -> None: + groups = self.create_unprocessed_groups(4) + group_ids = sorted(group.id for group in groups) + for group_id in group_ids: + process_group_log(group_id) + + with ( + override_options({"issues.derived.check-sample-size": 2}), + patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[1]), + patch.object(check_fresh_derived_data_batch, "delay") as mock_check, + ): + heal_stale_derived_data() + + mock_check.assert_called_once_with( + group_id_start=group_ids[1], + group_id_end=group_ids[2] + 1, + ) def test_respects_killswitch(self) -> None: groups = self.create_unprocessed_groups(1) @@ -339,6 +459,142 @@ def test_respects_max_tasks(self) -> None: assert mock_delay.call_count == 2 +@with_feature("projects:issue-action-log-write-to-db") +class CheckFreshDerivedDataBatchTest(DerivedDataTaskTestBase): + def test_checks_only_fresh_rows_inline(self) -> None: + groups = self.create_unprocessed_groups(3) + group_ids = sorted(group.id for group in groups) + for group_id in group_ids: + process_group_log(group_id) + GroupDerivedData.objects.filter(group_id=group_ids[1]).update(pipeline_hash="stale") + + with patch("sentry.issues.derived.tasks_util.metrics.incr") as mock_incr: + check_fresh_derived_data_batch( + group_id_start=group_ids[0], + group_id_end=group_ids[-1] + 1, + ) + + assert mock_incr.call_args_list == [ + call("issues.derived.check_group", sample_rate=1.0, tags={"result": "success"}), + call("issues.derived.check_group", sample_rate=1.0, tags={"result": "success"}), + ] + + def test_reschedules_timed_out_group_with_check_id(self) -> None: + group = self.create_unprocessed_groups(1)[0] + derived = process_group_log(group.id) + assert derived.pipeline_hash is not None + check_id = CheckId( + "invocation-id", + group.id, + derived.generated_at, + derived.cursor_date, + derived.cursor_id, + derived.pipeline_hash, + ) + + with ( + patch( + "sentry.issues.derived.check.check_derived_data", + side_effect=CheckTimeout(check_id), + ), + patch.object(check_fresh_derived_data_batch, "delay") as mock_delay, + ): + check_fresh_derived_data_batch( + group_id_start=group.id, + group_id_end=group.id + 1, + ) + + mock_delay.assert_called_once_with( + group_id_start=group.id, + group_id_end=group.id + 1, + resume_check_id="invocation-id", + resume_generated_at=derived.generated_at.isoformat(), + resume_cursor_date=derived.cursor_date.isoformat(), + resume_cursor_id=derived.cursor_id, + resume_pipeline_hash=derived.pipeline_hash, + prior_runs=1, + ) + + def test_advances_after_check_retry_limit(self) -> None: + group = self.create_unprocessed_groups(1)[0] + derived = process_group_log(group.id) + assert derived.pipeline_hash is not None + check_id = CheckId( + "invocation-id", + group.id, + derived.generated_at, + derived.cursor_date, + derived.cursor_id, + derived.pipeline_hash, + ) + + with ( + patch( + "sentry.issues.derived.check.check_derived_data", + side_effect=CheckTimeout(check_id), + ), + patch("sentry.issues.derived.tasks._MAX_CHECK_RUNS", 1), + patch.object(check_fresh_derived_data_batch, "delay") as mock_delay, + patch("sentry.issues.derived.tasks_util.metrics.incr") as mock_incr, + ): + check_fresh_derived_data_batch( + group_id_start=group.id, + group_id_end=group.id + 2, + ) + + mock_delay.assert_called_once_with( + group_id_start=group.id + 1, + group_id_end=group.id + 2, + ) + mock_incr.assert_called_once_with( + "issues.derived.check_group", + sample_rate=1.0, + tags={"result": "no_result"}, + ) + + +@with_feature("projects:issue-action-log-write-to-db") +class PickRandomFreshGroupRangeTest(DerivedDataTaskTestBase): + def test_returns_range_with_up_to_target_size(self) -> None: + groups = self.create_unprocessed_groups(4) + group_ids = sorted(group.id for group in groups) + for group_id in group_ids: + process_group_log(group_id) + + with patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[1]): + result = _pick_random_fresh_group_range(PIPELINE.pipeline_hash, target_size=2) + + assert result == (group_ids[1], group_ids[2] + 1) + + def test_returns_shorter_range_near_upper_bound(self) -> None: + groups = self.create_unprocessed_groups(3) + group_ids = sorted(group.id for group in groups) + for group_id in group_ids: + process_group_log(group_id) + + with patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[-1]): + result = _pick_random_fresh_group_range(PIPELINE.pipeline_hash, target_size=2) + + assert result == (group_ids[-1], group_ids[-1] + 1) + + def test_returns_none_without_fresh_rows(self) -> None: + assert _pick_random_fresh_group_range(PIPELINE.pipeline_hash, target_size=1000) is None + + def test_caps_target_size(self) -> None: + groups = self.create_unprocessed_groups(3) + group_ids = sorted(group.id for group in groups) + for group_id in group_ids: + process_group_log(group_id) + + with ( + patch("sentry.issues.derived.tasks_util._MAX_CHECK_GROUPS", 2), + patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[0]), + ): + result = _pick_random_fresh_group_range(PIPELINE.pipeline_hash, target_size=1000) + + assert result == (group_ids[0], group_ids[1] + 1) + + @with_feature("projects:issue-action-log-write-to-db") class RegenerateStaleDerivedDataBatchTest(DerivedDataTaskTestBase): @staticmethod From 657911b2b69014b37ef6e3ab946d4e4e92795073 Mon Sep 17 00:00:00 2001 From: Kyle Consalus Date: Fri, 7 Aug 2026 10:35:19 -0700 Subject: [PATCH 2/2] smarter --- src/sentry/issues/derived/tasks.py | 13 ++-- src/sentry/issues/derived/tasks_util.py | 29 ++++--- src/sentry/options/defaults.py | 6 +- .../sentry/issues/derived/test_processing.py | 3 + tests/sentry/issues/derived/test_tasks.py | 78 ++++++++++++++----- 5 files changed, 90 insertions(+), 39 deletions(-) diff --git a/src/sentry/issues/derived/tasks.py b/src/sentry/issues/derived/tasks.py index c2f1138d7007..f268b291f1f9 100644 --- a/src/sentry/issues/derived/tasks.py +++ b/src/sentry/issues/derived/tasks.py @@ -621,7 +621,7 @@ def heal_stale_derived_data(**kwargs: object) -> None: """Rebuild a chunk of GroupDerivedData rows whose ``pipeline_hash`` is stale/NULL.""" from sentry import options from sentry.issues.derived.processing import PIPELINE - from sentry.issues.derived.tasks_util import _pick_random_fresh_group_range + from sentry.issues.derived.tasks_util import _pick_random_fresh_group_ranges from sentry.issues.models.groupderiveddata import GroupDerivedData if not options.get("issues.derived.heal-enabled"): @@ -644,11 +644,12 @@ def heal_stale_derived_data(**kwargs: object) -> None: ) if not group_ids: logger.info("heal_stale_derived_data.nothing_to_heal") - check_range = _pick_random_fresh_group_range( - current_hash, options.get("issues.derived.check-sample-size") + check_ranges = _pick_random_fresh_group_ranges( + current_hash, + batch_size=batch_size, + task_count=options.get("issues.derived.check-task-count"), ) - if check_range is not None: - start, end = check_range + for start, end in check_ranges: check_fresh_derived_data_batch.delay( group_id_start=start, group_id_end=end, @@ -657,7 +658,7 @@ def heal_stale_derived_data(**kwargs: object) -> None: logger.info( "heal_stale_derived_data.checks_scheduled", extra={ - "scheduled": check_range is not None, + "task_count": len(check_ranges), "pipeline_hash": current_hash, }, ) diff --git a/src/sentry/issues/derived/tasks_util.py b/src/sentry/issues/derived/tasks_util.py index 3be3585bd9a1..93531dfc9eaf 100644 --- a/src/sentry/issues/derived/tasks_util.py +++ b/src/sentry/issues/derived/tasks_util.py @@ -33,28 +33,39 @@ def _record_check_result(result: CheckResult | None) -> None: ) -def _pick_random_fresh_group_range(pipeline_hash: str, target_size: int) -> tuple[int, int] | None: - """Pick a random range containing up to ``target_size`` fresh rows.""" - target_size = min(target_size, _MAX_CHECK_GROUPS) - if target_size <= 0: - return None +def _pick_random_fresh_group_ranges( + pipeline_hash: str, *, batch_size: int, task_count: int +) -> list[tuple[int, int]]: + """Pick contiguous check ranges from one random anchor (slide-to-fill if short).""" + if batch_size <= 0 or task_count <= 0: + return [] + need = min(batch_size * task_count, _MAX_CHECK_GROUPS) fresh = GroupDerivedData.objects.filter(pipeline_hash=pipeline_hash) bounds = fresh.aggregate(min_group_id=Min("group_id"), max_group_id=Max("group_id")) min_group_id = bounds["min_group_id"] max_group_id = bounds["max_group_id"] if min_group_id is None or max_group_id is None: - return None + return [] random_start = random.randint(min_group_id, max_group_id) group_ids = list( fresh.filter(group_id__gte=random_start) .order_by("group_id") - .values_list("group_id", flat=True)[:target_size] + .values_list("group_id", flat=True)[:need] ) + if len(group_ids) < need: + # Short forward tail: take the last ``need`` fresh rows (one contiguous band). + group_ids = list(fresh.order_by("-group_id").values_list("group_id", flat=True)[:need]) + group_ids.reverse() if not group_ids: - return None - return group_ids[0], group_ids[-1] + 1 + return [] + + ranges: list[tuple[int, int]] = [] + for i in range(0, len(group_ids), batch_size): + chunk = group_ids[i : i + batch_size] + ranges.append((chunk[0], chunk[-1] + 1)) + return ranges def _resume_check_id( diff --git a/src/sentry/options/defaults.py b/src/sentry/options/defaults.py index b5a93484ad63..16366df3413c 100644 --- a/src/sentry/options/defaults.py +++ b/src/sentry/options/defaults.py @@ -4113,10 +4113,10 @@ type=Int, flags=FLAG_AUTOMATOR_MODIFIABLE, ) -# Approximate number of fresh groups to check when there is no stale derived data to heal. +# Number of random check batches to schedule when there is no stale derived data to heal. register( - "issues.derived.check-sample-size", - default=1000, + "issues.derived.check-task-count", + default=5, type=Int, flags=FLAG_AUTOMATOR_MODIFIABLE, ) diff --git a/tests/sentry/issues/derived/test_processing.py b/tests/sentry/issues/derived/test_processing.py index 5f5ac80784b7..e9e1dd99e1e9 100644 --- a/tests/sentry/issues/derived/test_processing.py +++ b/tests/sentry/issues/derived/test_processing.py @@ -707,6 +707,8 @@ def test_pipeline_hash_null_stale_still_incrementally_updates(self) -> None: derived = process_group_log(group.id) first_cursor = derived.cursor_id + # Explicit date_added: db Now() is transaction-start time and can sort + # before a cursor stamped with Python timezone.now() from publish. new_entry = GroupActionLogEntry.objects.create( group_id=group.id, project_id=group.project_id, @@ -715,6 +717,7 @@ def test_pipeline_hash_null_stale_still_incrementally_updates(self) -> None: actor_id=0, source=SOURCE, data={}, + date_added=derived.cursor_date + timedelta(seconds=1), ) # Officially mark the row stale by resetting pipeline_hash to NULL. diff --git a/tests/sentry/issues/derived/test_tasks.py b/tests/sentry/issues/derived/test_tasks.py index 571ad8da4b74..49090e68bb88 100644 --- a/tests/sentry/issues/derived/test_tasks.py +++ b/tests/sentry/issues/derived/test_tasks.py @@ -16,7 +16,7 @@ heal_stale_derived_data, regenerate_stale_derived_data_batch, ) -from sentry.issues.derived.tasks_util import _pick_random_fresh_group_range +from sentry.issues.derived.tasks_util import _pick_random_fresh_group_ranges from sentry.issues.models.groupderiveddata import GroupDerivedData from sentry.models.group import Group from sentry.testutils.cases import TestCase @@ -361,28 +361,34 @@ def test_no_stale_data(self) -> None: heal_stale_derived_data() mock_regenerate.assert_not_called() + # One anchor + contiguous fan-out; 2 groups fit in a single default batch. mock_check.assert_called_once_with( group_id_start=group_ids[0], group_id_end=group_ids[-1] + 1, ) - def test_checks_random_configured_group_range(self) -> None: + def test_schedules_contiguous_ranges_from_one_anchor(self) -> None: groups = self.create_unprocessed_groups(4) group_ids = sorted(group.id for group in groups) for group_id in group_ids: process_group_log(group_id) with ( - override_options({"issues.derived.check-sample-size": 2}), - patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[1]), + override_options( + { + "issues.derived.check-task-count": 2, + "issues.derived.heal-batch-size": 2, + } + ), + patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[0]), patch.object(check_fresh_derived_data_batch, "delay") as mock_check, ): heal_stale_derived_data() - mock_check.assert_called_once_with( - group_id_start=group_ids[1], - group_id_end=group_ids[2] + 1, - ) + assert mock_check.call_args_list == [ + call(group_id_start=group_ids[0], group_id_end=group_ids[1] + 1), + call(group_id_start=group_ids[2], group_id_end=group_ids[3] + 1), + ] def test_respects_killswitch(self) -> None: groups = self.create_unprocessed_groups(1) @@ -554,33 +560,61 @@ def test_advances_after_check_retry_limit(self) -> None: @with_feature("projects:issue-action-log-write-to-db") -class PickRandomFreshGroupRangeTest(DerivedDataTaskTestBase): - def test_returns_range_with_up_to_target_size(self) -> None: - groups = self.create_unprocessed_groups(4) +class PickRandomFreshGroupRangesTest(DerivedDataTaskTestBase): + def test_returns_contiguous_ranges_from_anchor(self) -> None: + groups = self.create_unprocessed_groups(6) group_ids = sorted(group.id for group in groups) for group_id in group_ids: process_group_log(group_id) with patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[1]): - result = _pick_random_fresh_group_range(PIPELINE.pipeline_hash, target_size=2) + result = _pick_random_fresh_group_ranges( + PIPELINE.pipeline_hash, batch_size=2, task_count=2 + ) + + # need=4 and 5 rows remain at/after anchor → no slide. + assert result == [ + (group_ids[1], group_ids[2] + 1), + (group_ids[3], group_ids[4] + 1), + ] + + def test_slides_window_to_fill_near_upper_bound(self) -> None: + groups = self.create_unprocessed_groups(5) + group_ids = sorted(group.id for group in groups) + for group_id in group_ids: + process_group_log(group_id) + + with patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[-1]): + result = _pick_random_fresh_group_ranges( + PIPELINE.pipeline_hash, batch_size=2, task_count=1 + ) - assert result == (group_ids[1], group_ids[2] + 1) + # need=2 but only 1 row forward of the anchor → last 2 fresh rows. + assert result == [(group_ids[-2], group_ids[-1] + 1)] - def test_returns_shorter_range_near_upper_bound(self) -> None: + def test_slides_to_all_rows_when_table_smaller_than_need(self) -> None: groups = self.create_unprocessed_groups(3) group_ids = sorted(group.id for group in groups) for group_id in group_ids: process_group_log(group_id) with patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[-1]): - result = _pick_random_fresh_group_range(PIPELINE.pipeline_hash, target_size=2) + result = _pick_random_fresh_group_ranges( + PIPELINE.pipeline_hash, batch_size=2, task_count=2 + ) - assert result == (group_ids[-1], group_ids[-1] + 1) + assert result == [ + (group_ids[0], group_ids[1] + 1), + (group_ids[2], group_ids[2] + 1), + ] - def test_returns_none_without_fresh_rows(self) -> None: - assert _pick_random_fresh_group_range(PIPELINE.pipeline_hash, target_size=1000) is None + def test_returns_empty_without_fresh_rows(self) -> None: + assert ( + _pick_random_fresh_group_ranges(PIPELINE.pipeline_hash, batch_size=1000, task_count=5) + == [] + ) - def test_caps_target_size(self) -> None: + def test_caps_total_groups(self) -> None: groups = self.create_unprocessed_groups(3) group_ids = sorted(group.id for group in groups) for group_id in group_ids: @@ -590,9 +624,11 @@ def test_caps_target_size(self) -> None: patch("sentry.issues.derived.tasks_util._MAX_CHECK_GROUPS", 2), patch("sentry.issues.derived.tasks_util.random.randint", return_value=group_ids[0]), ): - result = _pick_random_fresh_group_range(PIPELINE.pipeline_hash, target_size=1000) + result = _pick_random_fresh_group_ranges( + PIPELINE.pipeline_hash, batch_size=1000, task_count=5 + ) - assert result == (group_ids[0], group_ids[1] + 1) + assert result == [(group_ids[0], group_ids[1] + 1)] @with_feature("projects:issue-action-log-write-to-db")