Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 26 additions & 6 deletions src/sentry/event_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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
Expand All @@ -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:
Expand Down
10 changes: 10 additions & 0 deletions src/sentry/options/defaults.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
32 changes: 29 additions & 3 deletions src/sentry/reprocessing2.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
Expand Down
9 changes: 9 additions & 0 deletions src/sentry/services/eventstore/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
115 changes: 114 additions & 1 deletion tests/sentry/test_reprocessing2.py
Original file line number Diff line number Diff line change
@@ -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):
Expand Down Expand Up @@ -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)
Loading