diff --git a/src/sentry/workflow_engine/caches/detector.py b/src/sentry/workflow_engine/caches/detector.py index e171c0834486..06b53ea6f1df 100644 --- a/src/sentry/workflow_engine/caches/detector.py +++ b/src/sentry/workflow_engine/caches/detector.py @@ -46,7 +46,7 @@ def _query_detectors(source_id: str, query_type: str) -> list[Detector]: data_sources__source_id=source_id, data_sources__type=query_type, ) - .select_related("workflow_condition_group") + .select_related("project__organization", "workflow_condition_group") .prefetch_related("workflow_condition_group__conditions") .distinct() .order_by("id") diff --git a/src/sentry/workflow_engine/processors/__init__.py b/src/sentry/workflow_engine/processors/__init__.py index eee12a669cd8..3ff04c62f0d2 100644 --- a/src/sentry/workflow_engine/processors/__init__.py +++ b/src/sentry/workflow_engine/processors/__init__.py @@ -2,6 +2,12 @@ "DataConditionEvaluation", "DataConditionGroupEvaluation", "DetectorEvaluation", + "ProcessDetectorsResult", ] -from .evaluations import DataConditionEvaluation, DataConditionGroupEvaluation, DetectorEvaluation +from .evaluations import ( + DataConditionEvaluation, + DataConditionGroupEvaluation, + DetectorEvaluation, + ProcessDetectorsResult, +) diff --git a/src/sentry/workflow_engine/processors/detector.py b/src/sentry/workflow_engine/processors/detector.py index 5ee24794d103..3f8878fd0c23 100644 --- a/src/sentry/workflow_engine/processors/detector.py +++ b/src/sentry/workflow_engine/processors/detector.py @@ -10,6 +10,7 @@ from sentry.issues.producer import PayloadType, produce_occurrence_to_kafka from sentry.models.activity import Activity from sentry.models.group import Group +from sentry.models.organization import Organization from sentry.services.eventstore.models import GroupEvent from sentry.utils import metrics from sentry.utils.cache import cache @@ -21,7 +22,8 @@ ) from sentry.workflow_engine.models import DataPacket, Detector from sentry.workflow_engine.models.detector_group import DetectorGroup -from sentry.workflow_engine.processors import DetectorEvaluation +from sentry.workflow_engine.processors import DetectorEvaluation, ProcessDetectorsResult +from sentry.workflow_engine.processors.evaluation_logging import emit_detector_evaluation_logs from sentry.workflow_engine.types import ( DetectorGroupKey, DetectorId, @@ -271,6 +273,16 @@ def create_issue_platform_payload(result: DetectorEvaluation, detector_type: str ) +def _get_detector_organization(detector: Detector) -> Organization: + if detector.project_id is not None: + return detector.linked_project.organization + + organization_id = detector.config.get("organization_id") + if not isinstance(organization_id, int): + raise ValueError("Organization-scoped detector is missing organization_id") + return Organization.objects.get_from_cache(id=organization_id) + + @trace def process_detectors[T]( data_packet: DataPacket[T], detectors: list[Detector] @@ -293,6 +305,17 @@ def process_detectors[T]( ): detector_results = handler.evaluate(data_packet) + emit_detector_evaluation_logs( + logger, + organization=_get_detector_organization(detector), + result=ProcessDetectorsResult( + detector_id=detector.id, + detector_type=detector.type, + project_id=detector.project_id, + evaluations=detector_results, + ), + ) + for result in detector_results.values(): logger_extra = { "detector": detector.id, diff --git a/src/sentry/workflow_engine/processors/evaluation_logging.py b/src/sentry/workflow_engine/processors/evaluation_logging.py index 945620d8d674..abcd4e692509 100644 --- a/src/sentry/workflow_engine/processors/evaluation_logging.py +++ b/src/sentry/workflow_engine/processors/evaluation_logging.py @@ -2,19 +2,65 @@ import random from logging import Logger -from typing import TYPE_CHECKING +from typing import TYPE_CHECKING, cast from sentry import features, options from sentry.utils.sdk import sdk_logger +from sentry.workflow_engine.processors.evaluations.detector import ProcessDetectorsResult from sentry.workflow_engine.processors.evaluations.workflow import ProcessWorkflowsResult if TYPE_CHECKING: from sentry.models.organization import Organization +DETECTOR_EVALUATION_LOG_PREFIX = "workflow_engine.process_detectors.evaluation" WORKFLOW_EVALUATION_LOG_PREFIX = "workflow_engine.process_workflows.evaluation" +def _should_emit_evaluation_logs(organization: Organization) -> bool: + if features.has("organizations:workflow-engine-log-evaluations", organization): + return True + sample_rate = cast(float, options.get("workflow_engine.evaluation_log_sample_rate")) + return random.random() < sample_rate + + +def _emit_evaluation_artifacts( + logger: Logger, + *, + organization_id: int, + artifacts: list[dict[str, object]], + log_prefix: str, +) -> None: + direct_to_sentry = options.get("workflow_engine.evaluation_logs_direct_to_sentry") + for artifact in artifacts: + artifact["organization_id"] = organization_id + + if direct_to_sentry: + sdk_logger.info(log_prefix, attributes=artifact) + else: + logger.info(log_prefix, extra=artifact) + + +def emit_detector_evaluation_logs( + logger: Logger, + *, + organization: Organization, + result: ProcessDetectorsResult, + log_prefix: str = DETECTOR_EVALUATION_LOG_PREFIX, +) -> bool: + """Sample a detector and emit one self-contained artifact per grouped evaluation.""" + if not _should_emit_evaluation_logs(organization): + return False + + _emit_evaluation_artifacts( + logger, + organization_id=organization.id, + artifacts=result.evaluation_artifacts(), + log_prefix=log_prefix, + ) + return True + + def emit_workflow_evaluation_logs( logger: Logger, *, @@ -23,25 +69,18 @@ def emit_workflow_evaluation_logs( log_prefix: str = WORKFLOW_EVALUATION_LOG_PREFIX, ) -> bool: """Sample a batch and emit one self-contained artifact per workflow evaluation.""" - should_log = features.has("organizations:workflow-engine-log-evaluations", organization) - if not should_log: - should_log = random.random() < options.get("workflow_engine.evaluation_log_sample_rate") - - if not should_log: + if not _should_emit_evaluation_logs(organization): return False - direct_to_sentry = options.get("workflow_engine.evaluation_logs_direct_to_sentry") artifacts = ( [evaluation.to_artifact() for evaluation in result.evaluations.values()] if result.evaluations else [result.to_artifact()] ) - for artifact in artifacts: - artifact["organization_id"] = organization.id - - if direct_to_sentry: - sdk_logger.info(log_prefix, attributes=artifact) - else: - logger.info(log_prefix, extra=artifact) - + _emit_evaluation_artifacts( + logger, + organization_id=organization.id, + artifacts=artifacts, + log_prefix=log_prefix, + ) return True diff --git a/src/sentry/workflow_engine/processors/evaluations/__init__.py b/src/sentry/workflow_engine/processors/evaluations/__init__.py index a5004233ae58..d2fee23a3f87 100644 --- a/src/sentry/workflow_engine/processors/evaluations/__init__.py +++ b/src/sentry/workflow_engine/processors/evaluations/__init__.py @@ -4,7 +4,9 @@ "DataConditionGroupEvaluation", "DetectorEvaluation", "DetectorEvaluationData", + "DetectorEvaluationOutcome", "DeferredWorkflowEvaluationResult", + "ProcessDetectorsResult", "ProcessWorkflowsResult", "WorkflowEvaluation", "WorkflowEvaluationData", @@ -13,7 +15,12 @@ from .condition import DataConditionEvaluation, DataConditionEvaluationException from .condition_group import DataConditionGroupEvaluation -from .detector import DetectorEvaluation, DetectorEvaluationData +from .detector import ( + DetectorEvaluation, + DetectorEvaluationData, + DetectorEvaluationOutcome, + ProcessDetectorsResult, +) from .workflow import ( DeferredWorkflowEvaluationResult, ProcessWorkflowsResult, diff --git a/src/sentry/workflow_engine/processors/evaluations/detector.py b/src/sentry/workflow_engine/processors/evaluations/detector.py index 5c202d72ae7f..4b7004297275 100644 --- a/src/sentry/workflow_engine/processors/evaluations/detector.py +++ b/src/sentry/workflow_engine/processors/evaluations/detector.py @@ -1,4 +1,5 @@ from dataclasses import dataclass +from enum import StrEnum from typing import Any, TypedDict from sentry.workflow_engine.types import DetectorGroupKey, DetectorPriorityLevel, DetectorResult @@ -13,6 +14,11 @@ class DetectorEvaluationData(TypedDict): event_data: dict[str, Any] | None # TODO - improve this typing, for now migrating +class DetectorEvaluationOutcome(StrEnum): + COMPLETED = "completed" + NO_RESULTS = "no_results" + + @dataclass(frozen=True, kw_only=True) class DetectorEvaluation( BaseWorkflowEngineEvaluation[ @@ -49,3 +55,35 @@ def artifact_fields(self) -> dict[str, Any]: "priority": self.priority.value, "trigger_group_evaluation": self.data["trigger_group_evaluation"].to_artifact(), } + + +@dataclass(frozen=True, kw_only=True) +class ProcessDetectorsResult: + detector_id: int + detector_type: str + project_id: int | None + evaluations: dict[DetectorGroupKey, DetectorEvaluation] + + @property + def outcome(self) -> DetectorEvaluationOutcome: + if self.evaluations: + return DetectorEvaluationOutcome.COMPLETED + return DetectorEvaluationOutcome.NO_RESULTS + + def to_artifact(self) -> dict[str, object]: + return { + "detector_id": self.detector_id, + "detector_type": self.detector_type, + "project_id": self.project_id, + "outcome": self.outcome, + } + + def evaluation_artifacts(self) -> list[dict[str, object]]: + detector_artifact = self.to_artifact() + if not self.evaluations: + return [detector_artifact] + + return [ + {**detector_artifact, **evaluation.to_artifact()} + for evaluation in self.evaluations.values() + ] diff --git a/tests/sentry/workflow_engine/processors/test_detector.py b/tests/sentry/workflow_engine/processors/test_detector.py index 74d7a27da098..e60e32770e26 100644 --- a/tests/sentry/workflow_engine/processors/test_detector.py +++ b/tests/sentry/workflow_engine/processors/test_detector.py @@ -18,6 +18,7 @@ from sentry.services.eventstore.models import GroupEvent from sentry.testutils.cases import TestCase from sentry.testutils.helpers.datetime import freeze_time +from sentry.testutils.helpers.features import Feature from sentry.testutils.helpers.options import override_options from sentry.testutils.pytest.fixtures import django_db_all from sentry.types.activity import ActivityType @@ -28,6 +29,7 @@ from sentry.workflow_engine.handlers.detector.stateful import get_redis_client from sentry.workflow_engine.models import DataPacket, Detector, DetectorState from sentry.workflow_engine.models.detector_group import DetectorGroup +from sentry.workflow_engine.processors import ProcessDetectorsResult from sentry.workflow_engine.processors.detector import ( EventDetectors, associate_new_group_with_detector, @@ -37,6 +39,8 @@ get_preferred_detector, process_detectors, ) +from sentry.workflow_engine.processors.evaluation_logging import emit_detector_evaluation_logs +from sentry.workflow_engine.processors.evaluations import DetectorEvaluationOutcome from sentry.workflow_engine.types import ( DetectorPriorityLevel, WorkflowEventData, @@ -99,6 +103,154 @@ def test(self) -> None: priority=DetectorPriorityLevel.HIGH, ) + def test_logs_canonical_evaluation_artifact(self) -> None: + detector = self.create_detector(type=self.handler_type.slug) + data_packet = self.build_data_packet(secret="do-not-log") + + with ( + Feature({"organizations:workflow-engine-log-evaluations": True}), + override_options({"workflow_engine.evaluation_logs_direct_to_sentry": False}), + mock.patch("sentry.workflow_engine.processors.detector.logger") as mock_logger, + ): + process_detectors(data_packet, [detector]) + + mock_logger.info.assert_called_once_with( + "workflow_engine.process_detectors.evaluation", + extra={ + "detector_id": detector.id, + "detector_type": detector.type, + "project_id": detector.project_id, + "outcome": DetectorEvaluationOutcome.COMPLETED, + "group_key": None, + "priority": DetectorPriorityLevel.HIGH.value, + "trigger_group_evaluation": { + "logic_type": "any", + "result": True, + "condition_evaluations": [], + "triggered": True, + "error": None, + }, + "triggered": True, + "error": None, + "organization_id": self.organization.id, + }, + ) + assert "do-not-log" not in str(mock_logger.info.call_args) + + def test_logs_detector_with_no_evaluation_results(self) -> None: + detector = self.create_detector(type=self.handler_type.slug) + handler = detector.detector_handler + assert handler is not None + + with ( + Feature({"organizations:workflow-engine-log-evaluations": True}), + override_options({"workflow_engine.evaluation_logs_direct_to_sentry": False}), + mock.patch.object(type(handler), "evaluate", return_value={}), + mock.patch("sentry.workflow_engine.processors.detector.logger") as mock_logger, + ): + assert process_detectors(self.build_data_packet(), [detector]) == [] + + mock_logger.info.assert_called_once_with( + "workflow_engine.process_detectors.evaluation", + extra={ + "detector_id": detector.id, + "detector_type": detector.type, + "project_id": detector.project_id, + "outcome": DetectorEvaluationOutcome.NO_RESULTS, + "organization_id": self.organization.id, + }, + ) + + def test_detector_emitter_samples_once_for_grouped_results(self) -> None: + detector, _ = self.create_detector_and_condition(type=self.handler_state_type.slug) + handler = detector.detector_handler + assert handler is not None + evaluations = handler.evaluate( + DataPacket("1", {"dedupe": 2, "group_vals": {"group_1": 6, "group_2": 10}}) + ) + mock_logger = mock.MagicMock() + + with ( + Feature({"organizations:workflow-engine-log-evaluations": False}), + override_options( + { + "workflow_engine.evaluation_log_sample_rate": 0.5, + "workflow_engine.evaluation_logs_direct_to_sentry": False, + } + ), + mock.patch( + "sentry.workflow_engine.processors.evaluation_logging.random.random", + return_value=0.1, + ) as mock_random, + ): + assert emit_detector_evaluation_logs( + mock_logger, + organization=self.organization, + result=ProcessDetectorsResult( + detector_id=detector.id, + detector_type=detector.type, + project_id=detector.project_id, + evaluations=evaluations, + ), + ) + + mock_random.assert_called_once_with() + assert mock_logger.info.call_count == 2 + assert {item.kwargs["extra"]["group_key"] for item in mock_logger.info.call_args_list} == { + "group_1", + "group_2", + } + + def test_detector_emitter_can_log_directly_to_sentry(self) -> None: + detector = self.create_detector(type=self.handler_type.slug) + handler = detector.detector_handler + assert handler is not None + evaluations = handler.evaluate(self.build_data_packet()) + mock_logger = mock.MagicMock() + + with ( + Feature({"organizations:workflow-engine-log-evaluations": True}), + override_options({"workflow_engine.evaluation_logs_direct_to_sentry": True}), + mock.patch( + "sentry.workflow_engine.processors.evaluation_logging.sdk_logger" + ) as mock_sentry_logger, + ): + assert emit_detector_evaluation_logs( + mock_logger, + organization=self.organization, + result=ProcessDetectorsResult( + detector_id=detector.id, + detector_type=detector.type, + project_id=detector.project_id, + evaluations=evaluations, + ), + ) + + mock_sentry_logger.info.assert_called_once_with( + "workflow_engine.process_detectors.evaluation", + attributes={ + **ProcessDetectorsResult( + detector_id=detector.id, + detector_type=detector.type, + project_id=detector.project_id, + evaluations=evaluations, + ).evaluation_artifacts()[0], + "organization_id": self.organization.id, + }, + ) + mock_logger.info.assert_not_called() + + def test_all_projects_detector_uses_configured_organization(self) -> None: + detector = self.create_detector(type=self.handler_type.slug) + detector.update(project=None, config={"organization_id": self.organization.id}) + + with mock.patch( + "sentry.workflow_engine.processors.detector.emit_detector_evaluation_logs" + ) as mock_emit: + process_detectors(self.build_data_packet(), [detector]) + + assert mock_emit.call_args.kwargs["organization"] == self.organization + @mock.patch("sentry.workflow_engine.processors.detector.produce_occurrence_to_kafka") def test_state_results(self, mock_produce_occurrence_to_kafka: MagicMock) -> None: detector, _ = self.create_detector_and_condition(type=self.handler_state_type.slug) @@ -293,8 +445,7 @@ def test_metrics_and_logs_fire( ), ], ) - assert mock_logger.info.call_count == 1 - assert mock_logger.info.call_args[0][0] == "detector_triggered" + assert any(call.args[0] == "detector_triggered" for call in mock_logger.info.call_args_list) @mock.patch("sentry.workflow_engine.processors.detector.produce_occurrence_to_kafka") @mock.patch("sentry.workflow_engine.processors.detector.metrics") @@ -341,8 +492,7 @@ def test_metrics_and_logs_resolve( ), ], ) - assert mock_logger.info.call_count == 2 - assert mock_logger.info.call_args[0][0] == "detector_resolved" + assert any(call.args[0] == "detector_resolved" for call in mock_logger.info.call_args_list) def test_doesnt_send_metric(self) -> None: detector = self.create_detector(type=self.no_handler_type.slug)