diff --git a/src/sentry/event_manager.py b/src/sentry/event_manager.py index 11186d425850..1de31a921317 100644 --- a/src/sentry/event_manager.py +++ b/src/sentry/event_manager.py @@ -125,7 +125,11 @@ 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_backup, + get_unprocessed_backup, + 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 ( @@ -1087,14 +1091,25 @@ 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. Depending on the status of + # `store.reprocessing-nodestore-backup.rollout`, the unprocessed copy is parked + # in either store; whichever holds it, it gets persisted here as a subkey. + from_node = False 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}), - unprocessed=True, - ) + unprocessed = get_unprocessed_backup(event.project_id, event.event_id) + from_node = unprocessed is not None + if unprocessed is None: + unprocessed = event_processing_store.get( + cache_key_for_event({"project": event.project_id, "event_id": event.event_id}), + unprocessed=True, + ) if unprocessed is not None: subkeys["unprocessed"] = unprocessed + metrics.incr( + "events.unprocessed_copy.promoted", + tags={"source": "node" if from_node else "processing_store"}, + sample_rate=1.0, + ) if app_feature: event_size = 0 @@ -1110,6 +1125,11 @@ def _nodestore_save_many(jobs: Sequence[Job], app_feature: str) -> None: job["event"].data["nodestore_insert"] = inserted_time job["event"].data.save(subkeys=subkeys) + # Only now that the payload is durable as a subkey is it safe to drop it from + # the temporary store. + if from_node: + delete_unprocessed_backup(event.project_id, event.event_id) + def _eventstream_insert_many(jobs: Sequence[Job]) -> None: for job in jobs: diff --git a/src/sentry/options/defaults.py b/src/sentry/options/defaults.py index 458f2d04a2cf..bb92dd6ff027 100644 --- a/src/sentry/options/defaults.py +++ b/src/sentry/options/defaults.py @@ -1121,6 +1121,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 to Nodestore instead of Redis +# for temporary holding until `event_manager` persists it as a subkey on the final +# event. +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..951a9dc62914 100644 --- a/src/sentry/reprocessing2.py +++ b/src/sentry/reprocessing2.py @@ -143,6 +143,11 @@ # and after which we just give up and mark the group as finished. REPROCESSING_TIMEOUT = 20 * 60 +# Unprocessed event copies are ultimately stored as subkeys in the same row as the +# processed event copy. While processing is still ongoing, they are stored in their own +# nodestore row with a 24h TTL. +UNPROCESSED_COPY_TTL = timedelta(hours=24) + # Note: This list of reasons is exposed in the EventReprocessableEndpoint to # the frontend. @@ -164,14 +169,35 @@ 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), ttl=UNPROCESSED_COPY_TTL) + else: + # Once the rollout is complete, this branch goes away. + event_processing_store.store(dict(data), unprocessed=True) + + +def get_unprocessed_backup(project_id: int, event_id: str) -> Any | None: + """ + Read the short-lived backup written by `backup_unprocessed_event`, not the durable + `unprocessed` subkey of a saved event. + """ + return nodestore.backend.get(Event.generate_unprocessed_node_id(project_id, event_id)) + + +def delete_unprocessed_backup(project_id: int, event_id: str) -> None: + """ + Drop the short-lived backup once its payload is durable elsewhere. Backups belonging + to events that never get saved are left to expire via `UNPROCESSED_COPY_TTL`. + """ + nodestore.backend.delete(Event.generate_unprocessed_node_id(project_id, event_id)) @dataclass 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/tests/sentry/test_reprocessing2.py b/tests/sentry/test_reprocessing2.py index a2f536314a7f..db528b5fc481 100644 --- a/tests/sentry/test_reprocessing2.py +++ b/tests/sentry/test_reprocessing2.py @@ -1,11 +1,26 @@ from __future__ import annotations +from typing import Any 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 ( + UNPROCESSED_COPY_TTL, + CannotReprocess, + _maybe_copy_attachment_into_cache, + backup_unprocessed_event, + delete_unprocessed_backup, + get_unprocessed_backup, + 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 +52,101 @@ 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[str, Any]: + 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_backup(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_backup(self.project.id, self.event_id) + + assert nodestore.backend.get(self._node_id()) is None + + @override_options({"store.reprocessing-nodestore-backup.rollout": 1.0}) + @mock.patch("sentry.reprocessing2.nodestore.backend.set") + def test_backup_expires_before_the_event_does(self, mock_set: mock.Mock) -> None: + backup_unprocessed_event(self._payload()) + + assert mock_set.call_args.kwargs["ttl"] == UNPROCESSED_COPY_TTL + + @override_options({"store.reprocessing-nodestore-backup.rollout": 1.0}) + def test_get_returns_payload_without_removing_node(self) -> None: + data = self._payload() + backup_unprocessed_event(data) + + assert get_unprocessed_backup(self.project.id, self.event_id) == data + assert nodestore.backend.get(self._node_id()) == data + + def test_get_without_node_returns_none(self) -> None: + assert get_unprocessed_backup(self.project.id, self.event_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 + ) + + def test_reads_subkey(self) -> None: + nodestore.backend.set_subkeys( + Event.generate_node_id(self.project.id, self.event_id), + {None: {"message": "processed"}, "unprocessed": self.unprocessed}, + ) + + assert pull_event_data(self.project.id, self.event_id).data == self.unprocessed + + def test_ignores_dedicated_node(self) -> None: + nodestore.backend.set( + Event.generate_unprocessed_node_id(self.project.id, self.event_id), self.unprocessed + ) + + with pytest.raises(CannotReprocess): + pull_event_data(self.project.id, self.event_id) + + def test_raises_when_no_copy_exists(self) -> None: + with pytest.raises(CannotReprocess): + pull_event_data(self.project.id, self.event_id)