Skip to content

Commit a17c071

Browse files
authored
fix(issue detectors): Fix missing nodestore data on segment-derived occurrences (#123288)
NOTE: This is heavily based on the work done in #123147, but includes a number of additional fixes identified by Claude as needing to be part of that change. Because the final result differs quite a bit from that PR, I created a new one which supersedes it. H/t to Claude for a great deal of help with the reasoning and rough drafts of all of the code included here, as well as all the tests. ------------------ Currently, occurrences from the span segment pipeline set `is_buffered_spans`, which builds an `Event` backed only by `snuba_data` and never writes to nodestore. The eventstream payload is therefore just `{"received": ...}`, so occurrences are reaching Snuba and showing up in the issue feed, but issue details is broken because there's no corresponding event JSON. To fix this problem, this PR drops the special case and sends segment-derived occurrences through the same path that every other occurrence takes. Each detected problem now carries its own event holding only the spans its evidence points at, filtered down to the fields issue details and Seer actually read, and trimmed to safely fit within the occurrence producer's limits. Notes: - Since the `event` we create is entirely synthetic, the occurrence id is reused for the event id. - Trimming happens both in terms of how many spans we keep and in terms of the data in those spans. Offender spans are given priority over parent and cause spans, since offender spans are the ones actually at the root of whatever problem we're reporting. - Any span whose id is referenced in one of the span id lists (parent, cause, or offender) but which isn't present in the segment is dropped, so we never error out trying to pull data from a span we don't have. - This fixes new occurrences only. Ones already written by the old path still point at nodestore keys that were never created, and stay broken until they age out. - There are a number of other issues Claude noted when working on this, but for ease of review those have been split off into follow-up PRs.
1 parent 54b1917 commit a17c071

5 files changed

Lines changed: 317 additions & 121 deletions

File tree

src/sentry/issues/occurrence_consumer.py

Lines changed: 3 additions & 55 deletions
Original file line numberDiff line numberDiff line change
@@ -109,51 +109,6 @@ def lookup_event(project_id: int, event_id: str) -> Event:
109109
return event
110110

111111

112-
@trace
113-
def create_event(project_id: int, event_id: str, event_data: dict[str, Any]) -> Event:
114-
return Event(
115-
event_id=event_id,
116-
project_id=project_id,
117-
# `snuba_data` only backs the `Event` property accessors. The eventstream serializes
118-
# `event.data`, where the `generic-events` schema requires `received`.
119-
data={"received": event_data["received"]},
120-
snuba_data={
121-
"event_id": event_data["event_id"],
122-
"project_id": event_data["project_id"],
123-
"timestamp": event_data["timestamp"],
124-
"release": event_data.get("release"),
125-
"environment": event_data.get("environment"),
126-
"platform": event_data.get("platform"),
127-
"tags.key": [tag[0] for tag in event_data.get("tags") or []],
128-
"tags.value": [tag[1] for tag in event_data.get("tags") or []],
129-
},
130-
)
131-
132-
133-
@trace
134-
def create_event_and_issue_occurrence(
135-
occurrence_data: IssueOccurrenceData, event_data: dict[str, Any]
136-
) -> tuple[IssueOccurrence, GroupInfo | None]:
137-
"""With standalone span ingestion, we won't be storing events in
138-
nodestore, so instead we create a light-weight event with a small
139-
set of fields that lets us create occurrences.
140-
"""
141-
project_id = occurrence_data["project_id"]
142-
event_id = occurrence_data["event_id"]
143-
if occurrence_data["event_id"] != event_data["event_id"]:
144-
raise ValueError(
145-
f"event_id in occurrence({occurrence_data['event_id']}) is different from event_id in event_data({event_data['event_id']})"
146-
)
147-
148-
event = create_event(project_id, event_id, event_data)
149-
150-
with metrics.timer(
151-
"occurrence_consumer._process_message.save_issue_occurrence",
152-
tags={"method": "create_event_and_issue_occurrence"},
153-
):
154-
return save_issue_occurrence(occurrence_data, event)
155-
156-
157112
@trace
158113
def process_event_and_issue_occurrence(
159114
occurrence_data: IssueOccurrenceData, event_data: dict[str, Any]
@@ -278,6 +233,7 @@ def _get_kwargs(payload: Mapping[str, Any]) -> Mapping[str, Any]:
278233
"request",
279234
"sdk",
280235
"server_name",
236+
"spans",
281237
"stacktrace",
282238
"trace_id",
283239
"transaction",
@@ -335,11 +291,7 @@ def _get_kwargs(payload: Mapping[str, Any]) -> Mapping[str, Any]:
335291
"title": occurrence_data["issue_title"],
336292
}
337293

338-
return {
339-
"occurrence_data": occurrence_data,
340-
"event_data": event_data,
341-
"is_buffered_spans": payload.get("is_buffered_spans") is True,
342-
}
294+
return {"occurrence_data": occurrence_data, "event_data": event_data}
343295
else:
344296
if not payload.get("event_id"):
345297
raise InvalidEventPayloadError(
@@ -363,8 +315,6 @@ def process_occurrence_message(
363315
kwargs = _get_kwargs(message)
364316
occurrence_data = kwargs["occurrence_data"]
365317
metric_tags = {"occurrence_type": occurrence_data["type"]}
366-
is_buffered_spans = kwargs.get("is_buffered_spans", False)
367-
368318
metrics.incr(
369319
"occurrence_ingest.messages",
370320
sample_rate=1.0,
@@ -399,9 +349,7 @@ def process_occurrence_message(
399349
set_span_tag(span, "result", "dropped_rate_limited")
400350
return None
401351

402-
if "event_data" in kwargs and is_buffered_spans:
403-
return create_event_and_issue_occurrence(kwargs["occurrence_data"], kwargs["event_data"])
404-
elif "event_data" in kwargs:
352+
if "event_data" in kwargs:
405353
set_span_tag(span, "result", "success")
406354
with metrics.timer(
407355
"occurrence_consumer._process_message.process_event_and_issue_occurrence",

src/sentry/issues/producer.py

Lines changed: 1 addition & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -55,10 +55,9 @@ def produce_occurrence_to_kafka(
5555
occurrence: IssueOccurrence | None = None,
5656
status_change: StatusChangeMessage | None = None,
5757
event_data: dict[str, Any] | None = None,
58-
is_buffered_spans: bool | None = False,
5958
) -> None:
6059
if payload_type == PayloadType.OCCURRENCE:
61-
payload_data = _prepare_occurrence_message(occurrence, event_data, is_buffered_spans)
60+
payload_data = _prepare_occurrence_message(occurrence, event_data)
6261
elif payload_type == PayloadType.STATUS_CHANGE:
6362
payload_data = _prepare_status_change_message(status_change)
6463
else:
@@ -96,7 +95,6 @@ def produce_occurrence_to_kafka(
9695
def _prepare_occurrence_message(
9796
occurrence: IssueOccurrence | None,
9897
event_data: dict[str, Any] | None,
99-
is_buffered_spans: bool | None = False,
10098
) -> MutableMapping[str, Any] | None:
10199
if not occurrence:
102100
raise ValueError("occurrence must be provided")
@@ -131,9 +129,6 @@ def _prepare_occurrence_message(
131129

132130
payload_data["event"] = event_data
133131

134-
if is_buffered_spans:
135-
payload_data["is_buffered_spans"] = True
136-
137132
return payload_data
138133

139134

src/sentry/spans/consumers/process_segments/message.py

Lines changed: 153 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -34,7 +34,7 @@
3434
get_detection_settings,
3535
)
3636
from sentry.issue_detection.performance_problem import PerformanceProblem
37-
from sentry.issues.issue_occurrence import IssueOccurrence
37+
from sentry.issues.issue_occurrence import IssueEvidence, IssueOccurrence
3838
from sentry.issues.producer import PayloadType, produce_occurrence_to_kafka
3939
from sentry.killswitches import killswitch_matches_context
4040
from sentry.models.environment import Environment
@@ -62,6 +62,31 @@
6262

6363
outcome_aggregator = OutcomeAggregator()
6464

65+
# This set of constants helps us keep occurrence data within the occurrence consumer's limits.
66+
# Limits on span count are split into overall and context (parent and cause) span limits in order to
67+
# prioritize offender spans, because they're the ones containing the actual problems we report.
68+
OVERALL_MAX_EVIDENCE_SPANS = 100
69+
MAX_EVIDENCE_CONTEXT_SPANS = 10
70+
MAX_SPAN_DESCRIPTION_LENGTH = 2048
71+
MAX_SPAN_DATA_VALUE_LENGTH = 500
72+
MAX_EVIDENCE_VALUE_LENGTH = 500
73+
MAX_EVIDENCE_LIST_ITEMS = 100
74+
75+
# The `evidence_data` values we actually use for issue details and Seer - all others are dropped
76+
# when we produce the occurrence
77+
EVIDENCE_SPAN_DATA_KEYS = frozenset(
78+
(
79+
"code.filepath",
80+
"code.function",
81+
"code.lineno",
82+
"http.query",
83+
"http.request.request_start",
84+
"http.request.response_start",
85+
"http.response_content_length",
86+
"url",
87+
)
88+
)
89+
6590

6691
@metrics.wraps("spans.consumers.process_segments.process_segment")
6792
def process_segment(
@@ -422,39 +447,153 @@ def _run_legacy_detectors(
422447
):
423448
return detected_problems
424449

425-
# Prepare a slimmer event payload for the occurrence consumer. This event will be persisted
426-
# by the consumer. Once issue detectors can run on standalone spans, we should directly
427-
# build a minimal occurrence event payload here, instead.
428-
event_data["spans"] = []
429-
event_data["timestamp"] = event_data["datetime"]
430-
450+
# Produce an occurrence for each problem, first filtering and trimming data to stay within the
451+
# occurrence consumer's limits
452+
spans_by_id = {span["span_id"]: span for span in event_data["spans"]}
431453
for problem in detected_problems:
454+
evidence_display = [
455+
IssueEvidence(
456+
evidence.name,
457+
_truncate_value_for_occurrence(evidence.value, MAX_EVIDENCE_VALUE_LENGTH),
458+
evidence.important,
459+
)
460+
for evidence in problem.evidence_display
461+
]
462+
evidence_data = _get_evidence_data_for_occurrence(problem, spans_by_id)
463+
464+
occurrence_id = uuid.uuid4().hex
465+
occurrence_spans = [
466+
_get_evidence_span_for_occurrence(spans_by_id[id])
467+
# We use `dict.fromkeys` here to preserve ordering
468+
for id in dict.fromkeys(
469+
(
470+
*evidence_data["parent_span_ids"],
471+
*evidence_data["cause_span_ids"],
472+
*evidence_data["offender_span_ids"],
473+
)
474+
)
475+
]
476+
occurrence_event_data = {
477+
**event_data,
478+
"event_id": occurrence_id,
479+
"spans": occurrence_spans,
480+
}
481+
432482
occurrence = IssueOccurrence(
433-
id=uuid.uuid4().hex,
483+
id=occurrence_id,
434484
resource_id=None,
435485
project_id=project.id,
436-
event_id=event_data["event_id"],
486+
event_id=occurrence_id,
437487
fingerprint=[problem.fingerprint],
438488
type=problem.type,
439489
issue_title=problem.title,
440-
subtitle=problem.desc,
490+
subtitle=_truncate_value_for_occurrence(problem.desc, MAX_EVIDENCE_VALUE_LENGTH),
441491
culprit=event_data["transaction"],
442-
evidence_data=problem.evidence_data or {},
443-
evidence_display=problem.evidence_display,
492+
evidence_data=evidence_data,
493+
evidence_display=evidence_display,
444494
detection_time=to_datetime(segment_span["end_timestamp"]),
445495
level="info",
446496
)
447497

448498
produce_occurrence_to_kafka(
449499
payload_type=PayloadType.OCCURRENCE,
450500
occurrence=occurrence,
451-
event_data=event_data,
452-
is_buffered_spans=True,
501+
event_data=occurrence_event_data,
453502
)
454503

455504
return detected_problems
456505

457506

507+
def _truncate_span_id_list(
508+
raw_span_ids: Sequence[str], spans_by_id: dict[str, Any], max_span_ids: int
509+
) -> list[str]:
510+
"""
511+
Drop ids the segment has no span for (so performance problem evidence never references a missing
512+
span), and then truncate the list of ids to the given max length.
513+
"""
514+
ids_for_existing_spans = [span_id for span_id in raw_span_ids if span_id in spans_by_id]
515+
return ids_for_existing_spans[:max_span_ids]
516+
517+
518+
def _get_evidence_data_for_occurrence(
519+
problem: PerformanceProblem, spans_by_id: dict[str, Any]
520+
) -> dict[str, Any]:
521+
# For the three lists of span ids, remove any which points to spans we don't have and then cap
522+
# the list length. The three lists can overlap, so the cap on offender span ids may be more
523+
# conservative than necessary, but better that than create an occurrence which gets rejected.
524+
parent_span_ids = _truncate_span_id_list(
525+
problem.parent_span_ids, spans_by_id, MAX_EVIDENCE_CONTEXT_SPANS
526+
)
527+
cause_span_ids = _truncate_span_id_list(
528+
problem.cause_span_ids, spans_by_id, MAX_EVIDENCE_CONTEXT_SPANS
529+
)
530+
max_offender_spans = OVERALL_MAX_EVIDENCE_SPANS - len(parent_span_ids) - len(cause_span_ids)
531+
offender_span_ids = _truncate_span_id_list(
532+
problem.offender_span_ids, spans_by_id, max_offender_spans
533+
)
534+
535+
# Now trim the rest of the evidence data
536+
evidence_data = problem.evidence_data or {}
537+
span_id_keys = {"parent_span_ids", "cause_span_ids", "offender_span_ids"}
538+
non_span_id_evidence_data = {
539+
key: value for key, value in evidence_data.items() if key not in span_id_keys
540+
}
541+
trimmed_non_span_id_evidence_data = _truncate_value_for_occurrence(
542+
non_span_id_evidence_data, MAX_EVIDENCE_VALUE_LENGTH
543+
)
544+
545+
return {
546+
**trimmed_non_span_id_evidence_data,
547+
"parent_span_ids": parent_span_ids,
548+
"cause_span_ids": cause_span_ids,
549+
"offender_span_ids": offender_span_ids,
550+
}
551+
552+
553+
def _truncate_value_for_occurrence(
554+
value: Any, max_chars: int, max_items: int = MAX_EVIDENCE_LIST_ITEMS
555+
) -> Any:
556+
"""
557+
Create a recursively-trimmed copy of a value for use in issue occurrences.
558+
"""
559+
if isinstance(value, str):
560+
return value[:max_chars]
561+
if isinstance(value, list | tuple):
562+
return [
563+
_truncate_value_for_occurrence(inner_value, max_chars, max_items)
564+
for inner_value in value[:max_items]
565+
]
566+
if isinstance(value, dict):
567+
return {
568+
key: _truncate_value_for_occurrence(inner_value, max_chars, max_items)
569+
for key, inner_value in value.items()
570+
}
571+
return value
572+
573+
574+
def _get_evidence_span_for_occurrence(span: dict[str, Any]) -> dict[str, Any]:
575+
trimmed_description = _truncate_value_for_occurrence(
576+
span.get("description"), MAX_SPAN_DESCRIPTION_LENGTH
577+
)
578+
# Only keep the `data` entries we actually need for evidence data
579+
filtered_data = {
580+
key: _truncate_value_for_occurrence(value, MAX_SPAN_DATA_VALUE_LENGTH)
581+
for key, value in (span.get("data") or {}).items()
582+
if key in EVIDENCE_SPAN_DATA_KEYS
583+
}
584+
585+
return {
586+
"span_id": span["span_id"],
587+
"trace_id": span.get("trace_id"),
588+
"op": span.get("op"),
589+
"description": trimmed_description,
590+
"start_timestamp": span.get("start_timestamp"),
591+
"timestamp": span.get("timestamp"),
592+
"exclusive_time": span.get("exclusive_time"),
593+
"data": filtered_data,
594+
}
595+
596+
458597
def _maybe_run_span_first_detector_parity_check(
459598
segment_span: CompatibleSpan,
460599
segment: list[CompatibleSpan],

tests/sentry/issues/test_occurrence_consumer.py

Lines changed: 0 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,6 @@
2424
)
2525
from sentry.issues.producer import _prepare_status_change_message
2626
from sentry.issues.status_change_message import StatusChangeMessage
27-
from sentry.models.environment import Environment
2827
from sentry.models.group import Group, GroupStatus
2928
from sentry.models.groupassignee import GroupAssignee
3029
from sentry.ratelimits.sliding_windows import Quota
@@ -310,51 +309,6 @@ def test_occurrence_rate_limit_quota(self) -> None:
310309
assert rate_limit_quota.limit == 1000
311310

312311

313-
class IssueOccurrenceBufferedSpansTest(IssueOccurrenceTestBase):
314-
"""Occurrences from the span segment processor never get saved to nodestore, so the eventstream
315-
payload has to be built entirely from the event data on the message.
316-
"""
317-
318-
def setUp(self) -> None:
319-
super().setUp()
320-
# Nothing on this path saves an event, so the default environment is never created for us
321-
Environment.get_or_create(self.project, None)
322-
323-
def run_buffered_spans_message(self, **event_overrides: Any) -> dict[str, Any]:
324-
message = get_test_message(self.project.id, is_buffered_spans=True)
325-
message["event"].update(event_overrides)
326-
327-
with mock.patch("sentry.issues.ingest.eventstream.backend.insert") as mock_insert:
328-
with self.feature("organizations:profile-file-io-main-thread-ingest"):
329-
result = _process_message(message)
330-
331-
assert result is not None
332-
assert mock_insert.call_count == 1
333-
group_event = mock_insert.call_args.kwargs["event"]
334-
return dict(group_event.get_raw_data(for_stream=True))
335-
336-
@django_db_all
337-
def test_received_reaches_the_eventstream(self) -> None:
338-
received = before_now(minutes=1)
339-
stream_data = self.run_buffered_spans_message(received=received.isoformat())
340-
341-
assert stream_data["received"] == pytest.approx(received.timestamp())
342-
343-
@django_db_all
344-
def test_received_reaches_the_eventstream_when_sent_as_a_timestamp(self) -> None:
345-
received = before_now(minutes=1).timestamp()
346-
stream_data = self.run_buffered_spans_message(received=received)
347-
348-
assert stream_data["received"] == pytest.approx(received)
349-
350-
@django_db_all
351-
@freeze_time()
352-
def test_null_received_falls_back_to_now(self) -> None:
353-
stream_data = self.run_buffered_spans_message(received=None)
354-
355-
assert stream_data["received"] == pytest.approx(datetime.datetime.now().timestamp())
356-
357-
358312
class IssueOccurrenceLookupEventIdTest(IssueOccurrenceTestBase):
359313
def test_lookup_event_doesnt_exist(self) -> None:
360314
message = get_test_message(self.project.id, include_event=False)

0 commit comments

Comments
 (0)