From 0a2c021ed1e1a7df8bb96eb9072a426d4ff5051c Mon Sep 17 00:00:00 2001 From: Matt Hammerly Date: Mon, 28 Sep 2026 18:29:16 -0700 Subject: [PATCH] ref(reprocessing): optionally write unprocessed payload to nodestore instead of redis --- src/sentry/deletions/tasks/nodestore.py | 24 ++--- src/sentry/event_manager.py | 8 +- src/sentry/options/defaults.py | 10 ++ src/sentry/reprocessing2.py | 38 +++++++- src/sentry/services/eventstore/models.py | 9 ++ src/sentry/tasks/reprocessing2.py | 5 +- src/sentry/tasks/store.py | 16 +++ src/sentry/tasks/symbolication.py | 7 ++ src/sentry/utils/cache.py | 17 +++- tests/sentry/test_reprocessing2.py | 119 ++++++++++++++++++++++- tests/sentry/utils/test_cache.py | 19 ++++ 11 files changed, 251 insertions(+), 21 deletions(-) create mode 100644 tests/sentry/utils/test_cache.py diff --git a/src/sentry/deletions/tasks/nodestore.py b/src/sentry/deletions/tasks/nodestore.py index 8600a0eb6b5c..5b1dcf88920b 100644 --- a/src/sentry/deletions/tasks/nodestore.py +++ b/src/sentry/deletions/tasks/nodestore.py @@ -247,17 +247,19 @@ def fetch_events_from_eventstore( def delete_events_from_nodestore(events: Sequence[Event], dataset: Dataset) -> None: - node_ids = [ - Event.generate_node_id( - event.project_id, - ( - event._snuba_data["occurrence_id"] - if dataset == Dataset.IssuePlatform - else event.event_id - ), - ) - for event in events - ] + node_ids = [] + for event in events: + if dataset == Dataset.IssuePlatform: + node_ids.append( + Event.generate_node_id(event.project_id, event._snuba_data["occurrence_id"]) + ) + else: + node_ids.append(Event.generate_node_id(event.project_id, event.event_id)) + # The unprocessed copy may be a separate node ratehr than a subkey. In this + # case, deleting the event doesn't take the unprocessed copy with it. + # N.B.: non-error events won't have one. + node_ids.append(Event.generate_unprocessed_node_id(event.project_id, event.event_id)) + nodestore.backend.delete_multi(node_ids) diff --git a/src/sentry/event_manager.py b/src/sentry/event_manager.py index 11186d425850..703ba3e80908 100644 --- a/src/sentry/event_manager.py +++ b/src/sentry/event_manager.py @@ -125,7 +125,7 @@ from sentry.receivers.features import record_event_processed from sentry.receivers.onboarding import record_release_received from sentry.releases.auto_creation import should_auto_create_releases -from sentry.reprocessing2 import is_reprocessed_event +from sentry.reprocessing2 import delete_unprocessed_event, is_reprocessed_event from sentry.seer.signed_seer_api import SeerViewerContext, make_signed_seer_api_request from sentry.services.eventstore.processing import event_processing_store from sentry.signals import ( @@ -555,6 +555,9 @@ def save_error_events( raise if not group_info: + # Returning here skips saving the event body in nodestore which renders the + # unprocessed copy unreachable. + delete_unprocessed_event(job["event"].project_id, job["event"].event_id) return job["event"] # store a reference to the group id to guarantee validation of isolation @@ -1087,7 +1090,8 @@ def _nodestore_save_many(jobs: Sequence[Job], app_feature: str) -> None: subkeys = {} event = job["event"] - # We only care about `unprocessed` for error events + # We only care about `unprocessed` for error events. Events whose unprocessed + # copy went straight to nodestore have nothing here and need no subkey. if event.get_event_type() not in ("transaction", "generic") and job["groups"]: unprocessed = event_processing_store.get( cache_key_for_event({"project": event.project_id, "event_id": event.event_id}), diff --git a/src/sentry/options/defaults.py b/src/sentry/options/defaults.py index 67be24864e43..58f16d964b73 100644 --- a/src/sentry/options/defaults.py +++ b/src/sentry/options/defaults.py @@ -1089,6 +1089,16 @@ # Killswitch to stop storing any reprocessing payloads. register("store.reprocessing-force-disable", default=False, flags=FLAG_AUTOMATOR_MODIFIABLE) +# Rollout for writing the unprocessed copy of an event straight to nodestore during +# preprocessing, instead of parking it in the processing store for event manager to +# promote to a subkey at save time. +register( + "store.reprocessing-nodestore-backup.rollout", + type=Float, + default=0.0, + flags=FLAG_MODIFIABLE_RATE | FLAG_AUTOMATOR_MODIFIABLE, +) + register( "store.ingest-events-raw-task.inline-save-event", type=Bool, diff --git a/src/sentry/reprocessing2.py b/src/sentry/reprocessing2.py index fc971ab1f2b6..1db52be34c95 100644 --- a/src/sentry/reprocessing2.py +++ b/src/sentry/reprocessing2.py @@ -164,14 +164,27 @@ def __init__(self, reason: CannotReprocessReason): def backup_unprocessed_event(data: Mapping[str, Any]) -> None: """ - Backup unprocessed event payload into redis. Only call if event should be - able to be reprocessed. + Backup unprocessed event payload. Only call if event should be able to be + reprocessed. """ if options.get("store.reprocessing-force-disable"): return - event_processing_store.store(dict(data), unprocessed=True) + if in_random_rollout("store.reprocessing-nodestore-backup.rollout"): + node_id = Event.generate_unprocessed_node_id(data["project"], data["event_id"]) + nodestore.backend.set(node_id, dict(data)) + else: + # Once the rollout is complete, this branch goes away along with the read in + # `_nodestore_save_many` that turns it into a subkey. + event_processing_store.store(dict(data), unprocessed=True) + + +def delete_unprocessed_event(project_id: int, event_id: str) -> None: + """ + Drop the unprocessed copy of an event that will never be saved. + """ + nodestore.backend.delete(Event.generate_unprocessed_node_id(project_id, event_id)) @dataclass @@ -191,8 +204,23 @@ def pull_event_data(project_id: int, event_id: str) -> ReprocessableEvent: raise CannotReprocess("event.not_found") with start_span(op="reprocess_events.nodestore.get", name="reprocess_events.nodestore.get"): - node_id = Event.generate_node_id(project_id, event_id) - data = nodestore.backend.get(node_id, subkey="unprocessed") + # If `store.reprocessing-nodestore-backup.rollout` was active for an event, it was + # saved in nodestore as a separate node. + data = nodestore.backend.get(Event.generate_unprocessed_node_id(project_id, event_id)) + source = "node" + if data is None: + # Fall back to reading the unprocessed copy from a subkey on the event node. + node_id = Event.generate_node_id(project_id, event_id) + data = nodestore.backend.get(node_id, subkey="unprocessed") + source = "subkey" if data is not None else "missing" + + # Tracks which representation served the read, so the rollout can be verified and + # the subkey fallback removed once it stops being hit. + metrics.incr( + "reprocessing2.unprocessed_copy.read", + tags={"source": source}, + sample_rate=1.0, + ) # Check data after checking presence of event to avoid too many instances. if data is None: diff --git a/src/sentry/services/eventstore/models.py b/src/sentry/services/eventstore/models.py index 3918bb27e0c9..46c2f5a88b1e 100644 --- a/src/sentry/services/eventstore/models.py +++ b/src/sentry/services/eventstore/models.py @@ -290,6 +290,15 @@ def generate_node_id(cls, project_id: int, event_id: str) -> str: """ return md5(f"{project_id}:{event_id}".encode()).hexdigest() + @classmethod + def generate_unprocessed_node_id(cls, project_id: int, event_id: str) -> str: + """ + Returns the node_id holding the unprocessed copy of an event, written before + symbolication so reprocessing can start over from the original payload. This + is a separate node from the event body, so it has to be deleted alongside it. + """ + return cls.generate_node_id(project_id, event_id) + ":u" + @property def project(self) -> Project: from sentry.models.project import Project diff --git a/src/sentry/tasks/reprocessing2.py b/src/sentry/tasks/reprocessing2.py index 6e69057b088f..8a5c0bccce34 100644 --- a/src/sentry/tasks/reprocessing2.py +++ b/src/sentry/tasks/reprocessing2.py @@ -266,8 +266,11 @@ def handle_remaining_events( for cls in EVENT_MODELS_TO_MIGRATE: cls.objects.filter(project_id=project_id, event_id__in=event_ids).delete() - # Remove from nodestore + # Remove from nodestore, including the separate unprocessed copies node_ids = [Event.generate_node_id(project_id, event_id) for event_id in event_ids] + node_ids += [ + Event.generate_unprocessed_node_id(project_id, event_id) for event_id in event_ids + ] nodestore.backend.delete_multi(node_ids) # Tell Snuba to delete the event data. diff --git a/src/sentry/tasks/store.py b/src/sentry/tasks/store.py index e5395f46967d..6ec4884b1960 100644 --- a/src/sentry/tasks/store.py +++ b/src/sentry/tasks/store.py @@ -33,6 +33,7 @@ issues_tasks, ) from sentry.utils import metrics +from sentry.utils.cache import event_from_cache_key from sentry.utils.event import track_event_since_received from sentry.utils.event_tracker import TransactionStageStatus, track_sampled_event from sentry.utils.safe import safe_execute @@ -358,6 +359,11 @@ def do_process_event( "events.failed", tags={"reason": "cache", "stage": "process"}, skip_internal=False ) error_logger.error("process.failed.empty", extra={"cache_key": cache_key}) + # The unprocessed copy was written during preprocessing, and the event will + # never reach nodestore now. + event_ref = event_from_cache_key(cache_key) + if event_ref is not None: + reprocessing2.delete_unprocessed_event(*event_ref) return track_event_since_received( @@ -579,6 +585,10 @@ def _do_save_event( metrics.incr( "events.failed", tags={"reason": "cache", "stage": "post"}, skip_internal=False ) + # The event will never reach nodestore, so its unprocessed copy is + # unreachable. This task does not retry. + if event_id and project_id: + reprocessing2.delete_unprocessed_event(project_id, event_id) return all_attachments = [] @@ -634,6 +644,12 @@ def _do_save_event( if cache_key: processing_store.delete_by_key(cache_key) + # Clean up the unprocessed copy from nodestore as it won't go through + # reprocessing. No-op if the unprocessed copy was stored in rc-processing as + # it will expire on its own. + if event_id: + reprocessing2.delete_unprocessed_event(project_id, event_id) + # Mark all the attachments as `rate_limited`, so they are being properly cleaned up in the `finally` block: for attachment in all_attachments: attachment.rate_limited = True diff --git a/src/sentry/tasks/symbolication.py b/src/sentry/tasks/symbolication.py index 418b2ceec548..67eaf7f8c254 100644 --- a/src/sentry/tasks/symbolication.py +++ b/src/sentry/tasks/symbolication.py @@ -6,6 +6,7 @@ import sentry_sdk from django.conf import settings +from sentry import reprocessing2 from sentry.killswitches import killswitch_matches_context from sentry.lang.native.processing import ( get_native_symbolication_functions, @@ -25,6 +26,7 @@ from sentry.tasks.base import instrumented_task from sentry.taskworker.namespaces import symbolication_tasks from sentry.utils import metrics +from sentry.utils.cache import event_from_cache_key from sentry.utils.sdk import set_current_event_project from sentry.utils.tracing import set_span_data, start_span @@ -75,6 +77,11 @@ def _do_symbolicate_event( "events.failed", tags={"reason": "cache", "stage": "symbolicate"}, skip_internal=False ) error_logger.error("symbolicate.failed.empty", extra={"cache_key": cache_key}) + # The unprocessed copy was written just before this task was queued, and the + # event will never reach nodestore now. + event_ref = event_from_cache_key(cache_key) + if event_ref is not None: + reprocessing2.delete_unprocessed_event(*event_ref) return event_id = str(data["event_id"]) diff --git a/src/sentry/utils/cache.py b/src/sentry/utils/cache.py index 8945ebd320dc..2d8b30f53a39 100644 --- a/src/sentry/utils/cache.py +++ b/src/sentry/utils/cache.py @@ -3,7 +3,7 @@ from django.core.cache import cache -__all__ = ["cache", "default_cache", "cache_key_for_event"] +__all__ = ["cache", "default_cache", "cache_key_for_event", "event_from_cache_key"] default_cache = cache @@ -12,3 +12,18 @@ def cache_key_for_event(data: Mapping[str, Any]) -> str: return "e:{}:{}".format(data["event_id"], data["project"]) + + +def event_from_cache_key(cache_key: str) -> tuple[int, str] | None: + """ + Inverse of `cache_key_for_event`, for pipeline steps that have lost the event body + and only hold its cache key. Returns `(project_id, event_id)`, or None if the key + is not in that shape. Never raises; callers are already on an error path. + """ + parts = cache_key.split(":") + if len(parts) != 3 or parts[0] != "e": + return None + try: + return int(parts[2]), parts[1] + except ValueError: + return None diff --git a/tests/sentry/test_reprocessing2.py b/tests/sentry/test_reprocessing2.py index a2f536314a7f..cfb7d9e912ac 100644 --- a/tests/sentry/test_reprocessing2.py +++ b/tests/sentry/test_reprocessing2.py @@ -2,10 +2,22 @@ from unittest import mock +import pytest + +from sentry import nodestore from sentry.models.eventattachment import EventAttachment -from sentry.reprocessing2 import _maybe_copy_attachment_into_cache +from sentry.reprocessing2 import ( + CannotReprocess, + _maybe_copy_attachment_into_cache, + backup_unprocessed_event, + delete_unprocessed_event, + pull_event_data, +) +from sentry.services.eventstore.models import Event +from sentry.services.eventstore.processing import event_processing_store from sentry.testutils.cases import TestCase from sentry.testutils.helpers.options import override_options +from sentry.utils.cache import cache_key_for_event class MaybeCopyAttachmentIntoCacheTest(TestCase): @@ -37,3 +49,108 @@ def test_objectstore_upload_stores_content_type(self, mock_get_session: mock.Moc assert cached.stored_id == "some-key" attachment.refresh_from_db() assert attachment.blob_path == "v2/some-key" + + +class UnprocessedCopyTest(TestCase): + event_id = "a" * 32 + + def _payload(self) -> dict: + return {"event_id": self.event_id, "project": self.project.id, "platform": "native"} + + def _node_id(self) -> str: + return Event.generate_unprocessed_node_id(self.project.id, self.event_id) + + @override_options({"store.reprocessing-nodestore-backup.rollout": 0.0}) + def test_backup_writes_to_processing_store(self) -> None: + data = self._payload() + backup_unprocessed_event(data) + + assert event_processing_store.get(cache_key_for_event(data), unprocessed=True) == data + assert nodestore.backend.get(self._node_id()) is None + + @override_options({"store.reprocessing-nodestore-backup.rollout": 1.0}) + def test_backup_writes_to_nodestore(self) -> None: + data = self._payload() + backup_unprocessed_event(data) + + assert nodestore.backend.get(self._node_id()) == data + assert event_processing_store.get(cache_key_for_event(data), unprocessed=True) is None + + def test_unprocessed_node_id_differs_from_event_node_id(self) -> None: + assert self._node_id() != Event.generate_node_id(self.project.id, self.event_id) + assert self._node_id().startswith(Event.generate_node_id(self.project.id, self.event_id)) + + @override_options({"store.reprocessing-nodestore-backup.rollout": 1.0}) + def test_delete_removes_node(self) -> None: + backup_unprocessed_event(self._payload()) + delete_unprocessed_event(self.project.id, self.event_id) + + assert nodestore.backend.get(self._node_id()) is None + + def test_delete_without_node_is_noop(self) -> None: + delete_unprocessed_event(self.project.id, self.event_id) + + assert nodestore.backend.get(self._node_id()) is None + + +class PullEventDataTest(TestCase): + event_id = "b" * 32 + + def setUp(self) -> None: + super().setUp() + self.unprocessed = { + "event_id": self.event_id, + "project": self.project.id, + "message": "unprocessed", + } + patcher = mock.patch("sentry.reprocessing2.eventstore.backend") + backend = patcher.start() + self.addCleanup(patcher.stop) + backend.get_event_by_id.return_value = Event( + project_id=self.project.id, event_id=self.event_id + ) + + @mock.patch("sentry.reprocessing2.metrics.incr") + def test_reads_separate_node(self, mock_incr: mock.Mock) -> None: + nodestore.backend.set( + Event.generate_unprocessed_node_id(self.project.id, self.event_id), self.unprocessed + ) + + result = pull_event_data(self.project.id, self.event_id) + + assert result.data == self.unprocessed + assert mock_incr.call_args.kwargs["tags"] == {"source": "node"} + + @mock.patch("sentry.reprocessing2.metrics.incr") + def test_falls_back_to_subkey(self, mock_incr: mock.Mock) -> None: + nodestore.backend.set_subkeys( + Event.generate_node_id(self.project.id, self.event_id), + {None: {"message": "processed"}, "unprocessed": self.unprocessed}, + ) + + result = pull_event_data(self.project.id, self.event_id) + + assert result.data == self.unprocessed + assert mock_incr.call_args.kwargs["tags"] == {"source": "subkey"} + + @mock.patch("sentry.reprocessing2.metrics.incr") + def test_prefers_separate_node_over_subkey(self, mock_incr: mock.Mock) -> None: + nodestore.backend.set( + Event.generate_unprocessed_node_id(self.project.id, self.event_id), self.unprocessed + ) + nodestore.backend.set_subkeys( + Event.generate_node_id(self.project.id, self.event_id), + {None: {"message": "processed"}, "unprocessed": {"message": "stale"}}, + ) + + result = pull_event_data(self.project.id, self.event_id) + + assert result.data == self.unprocessed + assert mock_incr.call_args.kwargs["tags"] == {"source": "node"} + + @mock.patch("sentry.reprocessing2.metrics.incr") + def test_raises_when_no_copy_exists(self, mock_incr: mock.Mock) -> None: + with pytest.raises(CannotReprocess): + pull_event_data(self.project.id, self.event_id) + + assert mock_incr.call_args.kwargs["tags"] == {"source": "missing"} diff --git a/tests/sentry/utils/test_cache.py b/tests/sentry/utils/test_cache.py new file mode 100644 index 000000000000..5790f3fd9cf5 --- /dev/null +++ b/tests/sentry/utils/test_cache.py @@ -0,0 +1,19 @@ +from __future__ import annotations + +import pytest + +from sentry.utils.cache import cache_key_for_event, event_from_cache_key + + +def test_round_trips_cache_key_for_event() -> None: + data = {"event_id": "a" * 32, "project": 42} + + assert event_from_cache_key(cache_key_for_event(data)) == (42, "a" * 32) + + +@pytest.mark.parametrize( + "cache_key", + ["", "e", "e:abc", "e:abc:123:456", "x:abc:123", "e:abc:notanint"], +) +def test_returns_none_for_unexpected_shapes(cache_key: str) -> None: + assert event_from_cache_key(cache_key) is None