Skip to content
Open
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
feat(celery): allow custom monitor config for beat tasks
Celery Beat auto-instrumentation derives the cron monitor config from the
schedule only, so max_runtime, checkin_margin, failure_issue_threshold,
recovery_threshold and owner cannot be set. Tasks whose runtime can exceed
Sentry's default 30-minute max_runtime are then reported as timed-out cron
failures even when they complete successfully.

Add a `beat_task_monitor_config` option to CeleryIntegration that maps beat
task names to partial monitor configs, merged over the derived config.

Refs #7838
  • Loading branch information
Alejandro
Alejandro committed Oct 2, 2026
commit 7e62674c6e903516171bda69e099ed9de046705c
17 changes: 14 additions & 3 deletions sentry_sdk/integrations/celery/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,15 @@
)

if TYPE_CHECKING:
from typing import Any, Callable, List, Optional, TypeVar, Union

from sentry_sdk._types import Event, EventProcessor, ExcInfo, Hint
from typing import Any, Callable, Dict, List, Optional, TypeVar, Union

from sentry_sdk._types import (
Event,
EventProcessor,
ExcInfo,
Hint,
MonitorConfig,
)

F = TypeVar("F", bound=Callable[..., Any])

Expand Down Expand Up @@ -63,10 +69,15 @@ def __init__(
propagate_traces: bool = True,
monitor_beat_tasks: bool = False,
exclude_beat_tasks: "Optional[List[str]]" = None,
beat_task_monitor_config: "Optional[Dict[str, MonitorConfig]]" = None,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Since we need a monitor_name provided as part of the configuration for any overrides to work, we can set the default value to an empty dictionary instead of None here.

Suggested change
beat_task_monitor_config: "Optional[Dict[str, MonitorConfig]]" = None,
beat_task_monitor_config: "Optional[Dict[str, MonitorConfig]]" = {},

This would allow us to update the (integration.beat_task_monitor_config or {}).get(monitor_name) conditional to either:

integration.beat_task_monitor_config.get(monitor_name)

or

getattr(integration.beat_task_monitor_config, monitor_name, None)

depending on your preference.

) -> None:
self.propagate_traces = propagate_traces
self.monitor_beat_tasks = monitor_beat_tasks
self.exclude_beat_tasks = exclude_beat_tasks
# Extra monitor config (e.g. ``max_runtime``, ``checkin_margin``) keyed
# by the name of the task in the Celery Beat schedule. The values are
# merged into the config the SDK derives from the schedule.
self.beat_task_monitor_config = beat_task_monitor_config

_patch_beat_apply_entry()
_patch_redbeat_apply_async()
Expand Down
15 changes: 13 additions & 2 deletions sentry_sdk/integrations/celery/beat.py
Original file line number Diff line number Diff line change
Expand Up @@ -57,7 +57,10 @@ def _get_headers(task: "Task") -> "dict[str, Any]":


def _get_monitor_config(
celery_schedule: "Any", app: "Celery", monitor_name: str
celery_schedule: "Any",
app: "Celery",
monitor_name: str,
monitor_config_overrides: "Optional[MonitorConfig]" = None,
) -> "MonitorConfig":
monitor_config: "MonitorConfig" = {}
schedule_type: "Optional[MonitorConfigScheduleType]" = None
Expand Down Expand Up @@ -111,6 +114,9 @@ def _get_monitor_config(
or "UTC"
)

if monitor_config_overrides is not None:
monitor_config.update(monitor_config_overrides)

return monitor_config


Expand All @@ -136,7 +142,12 @@ def _apply_crons_data_to_schedule_entry(
celery_schedule = schedule_entry.schedule
app = scheduler.app

monitor_config = _get_monitor_config(celery_schedule, app, monitor_name)
monitor_config = _get_monitor_config(
celery_schedule,
app,
monitor_name,
(integration.beat_task_monitor_config or {}).get(monitor_name),
)

is_supported_schedule = bool(monitor_config)
if not is_supported_schedule:
Expand Down
124 changes: 124 additions & 0 deletions tests/integrations/celery/test_celery_beat_crons.py
Original file line number Diff line number Diff line change
Expand Up @@ -385,6 +385,130 @@ def test_get_monitor_config_timezone_in_celery_schedule():
assert monitor_config["timezone"] == str(panama_tz)


def test_get_monitor_config_with_overrides():
app = MagicMock()
app.timezone = "Europe/Vienna"

celery_schedule = crontab(day_of_month="3", hour="12", minute="*/10")

monitor_config = _get_monitor_config(
celery_schedule,
app,
"foo",
{"max_runtime": 120, "checkin_margin": 5, "owner": "team:6"},
)

assert monitor_config == {
"schedule": {
"type": "crontab",
"value": "*/10 12 3 * *",
},
"timezone": "UTC",
"max_runtime": 120,
"checkin_margin": 5,
"owner": "team:6",
}


def test_get_monitor_config_overrides_can_replace_schedule():

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

We shouldn't support overriding the schedule (or its type) via this new property as there's validation that we do within the SDK to ensure the provided values are valid.

Instead, let's do the following:

1️⃣ Introduce a new type in _types.py to a subset of what's needed to support long-running tasks, etc. and update the relevant spots in this PR

MonitorConfigOverrides = TypedDict(
        "MonitorConfigOverrides",
        {
            "checkin_margin": int,
            "max_runtime": int,
            "failure_issue_threshold": int,
            "recovery_threshold": int,
            "owner": str,
        },
        total=False,
    )

2️⃣ Validate the overrides in the __init__ of the integration to ensure that everything's valid

In celery/beat.py:

_ALLOWED_MONITOR_CONFIG_OVERRIDE_KEYS = frozenset(
    {
        "checkin_margin",
        "max_runtime",
        "failure_issue_threshold",
        "recovery_threshold",
        "owner",
    }
)


def _validate_beat_task_monitor_config(
    beat_task_monitor_config: "Dict[str, MonitorConfigOverrides]",
) -> None:
    for task_name, overrides in beat_task_monitor_config.items():
        invalid = set(overrides) - _ALLOWED_MONITOR_CONFIG_OVERRIDE_KEYS
        if invalid:
            raise ValueError(
                f"Unsupported keys in beat_task_monitor_config for '{task_name}': "
                f"{sorted(invalid)}. Allowed keys: "
                f"{sorted(_ALLOWED_MONITOR_CONFIG_OVERRIDE_KEYS)}"
            )

and then invoke _validate_beat_task_monitor_config(beat_task_monitor_config) in the init code in __init__.py before assigning it to self.beat_task_monitor_config on CeleryIntegration.

app = MagicMock()
app.timezone = "Europe/Vienna"

celery_schedule = crontab(day_of_month="3", hour="12", minute="*/10")
override_schedule = {"type": "crontab", "value": "0 5 * * *"}

monitor_config = _get_monitor_config(
celery_schedule, app, "foo", {"schedule": override_schedule}
)

assert monitor_config["schedule"] == override_schedule


def test_beat_task_monitor_config_option():

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Let's remove this test

"""``beat_task_monitor_config`` is merged into the derived monitor config."""
fake_apply_entry = MagicMock()

fake_scheduler = MagicMock()
fake_scheduler.apply_entry = fake_apply_entry

fake_integration = MagicMock()
fake_integration.monitor_beat_tasks = True
fake_integration.exclude_beat_tasks = None
fake_integration.beat_task_monitor_config = {
"some_task_name": {"max_runtime": 120, "checkin_margin": 10}
}

fake_client = MagicMock()
fake_client.get_integration.return_value = fake_integration

fake_schedule_entry = MagicMock()
fake_schedule_entry.name = "some_task_name"
fake_schedule_entry.schedule = crontab(day_of_month="3", hour="12", minute="*/10")
fake_schedule_entry.options = {}

with mock.patch(
"sentry_sdk.integrations.celery.beat.Scheduler", fake_scheduler
) as Scheduler: # noqa: N806
with mock.patch(
"sentry_sdk.integrations.celery.sentry_sdk.get_client",
return_value=fake_client,
):
with mock.patch(
"sentry_sdk.integrations.celery.beat.capture_checkin",
return_value="check-in-id",
) as mock_capture_checkin:
# Mimic CeleryIntegration patching of Scheduler.apply_entry()
_patch_beat_apply_entry()
# Mimic Celery Beat calling a task from the Beat schedule
Scheduler.apply_entry(fake_scheduler, fake_schedule_entry)

assert fake_apply_entry.call_count == 1
monitor_config = mock_capture_checkin.call_args.kwargs["monitor_config"]
assert monitor_config["schedule"] == {
"type": "crontab",
"value": "*/10 12 3 * *",
}
assert monitor_config["max_runtime"] == 120
assert monitor_config["checkin_margin"] == 10


def test_beat_task_monitor_config_option_only_applies_to_matching_task():
fake_apply_entry = MagicMock()

fake_scheduler = MagicMock()
fake_scheduler.apply_entry = fake_apply_entry
Comment on lines +478 to +479

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Within the SDK, we're trying to move towards testing functionality via the public APIs rather than the private methods and use mocks as sparingly as we can.

Instead of this approach, can we look to instead do something similar to the test_beat_task_crons_success in test_celery_beat_cron_monitoring.py (which creates a celery app and asserts that the envelopes contain the right data)?


fake_integration = MagicMock()
fake_integration.monitor_beat_tasks = True
fake_integration.exclude_beat_tasks = None
fake_integration.beat_task_monitor_config = {"some_task_name": {"max_runtime": 120}}

fake_client = MagicMock()
fake_client.get_integration.return_value = fake_integration

fake_schedule_entry = MagicMock()
fake_schedule_entry.name = "another_task_name"
fake_schedule_entry.schedule = crontab(day_of_month="3", hour="12", minute="*/10")
fake_schedule_entry.options = {}

with mock.patch(
"sentry_sdk.integrations.celery.beat.Scheduler", fake_scheduler
) as Scheduler: # noqa: N806
with mock.patch(
"sentry_sdk.integrations.celery.sentry_sdk.get_client",
return_value=fake_client,
):
with mock.patch(
"sentry_sdk.integrations.celery.beat.capture_checkin",
return_value="check-in-id",
) as mock_capture_checkin:
_patch_beat_apply_entry()
Scheduler.apply_entry(fake_scheduler, fake_schedule_entry)

monitor_config = mock_capture_checkin.call_args.kwargs["monitor_config"]
assert "max_runtime" not in monitor_config


@pytest.mark.parametrize(
"task_name,exclude_beat_tasks,task_in_excluded_beat_tasks",
[
Expand Down