diff --git a/litellm/proxy/common_utils/scheduled_job_stagger.py b/litellm/proxy/common_utils/scheduled_job_stagger.py index e48e9686f13..05059cb6a6f 100644 --- a/litellm/proxy/common_utils/scheduled_job_stagger.py +++ b/litellm/proxy/common_utils/scheduled_job_stagger.py @@ -15,6 +15,13 @@ instant after a restart. Hashing rather than randomising keeps a given process's stable for its whole life and lets the applied offsets be logged once and reasoned about later. +A job that elects an owner is the exception: it drops ``identity`` and offsets by job id +alone. Only one replica does its work, so spreading it wins nothing, and spreading it +costs correctness, because a lease released when the body returns dedupes for that body's +runtime rather than for the lock's TTL. Replicas placed further apart than that each find +the lock free and each run. Sharing one instant per job keeps the burst apart job by job, +which is what this module is for, while letting the election do the rest. + The offset lives in the trigger rather than in a one-off ``next_run_time`` because a cron trigger recomputes each fire from the wall clock and would otherwise snap straight back onto the shared instant after its first shifted run. @@ -51,17 +58,19 @@ from litellm.constants import ( PTU_ROLLUP_LOCK_TTL_SECONDS, ) from litellm.proxy._types import ScheduledJobStaggerSettings +from litellm.proxy.common_utils.single_owner_job import SINGLE_OWNER_JOB_IDS GENERAL_SETTINGS_KEY: Final = "scheduled_job_stagger" #: Cron schedules LiteLLM picks on the operator's behalf, so shifting them changes nothing the #: operator asked for. Every other cron trigger is an operator-supplied crontab, preserved exactly. #: -#: The value is the span over which a second firing would redo work the first already did, which -#: is how long each job's leader-election lock stays held. Two replicas further apart than that -#: both find the key free and both run, which for the spend report means the customer gets it -#: twice. Offsets for these jobs are bounded by it, so widening the window cannot resurrect the -#: duplicate-work failure this feature exists to avoid. +#: The value bounds how far apart two replicas may be placed. It is a second line of defence +#: rather than the mechanism: an id in ``SINGLE_OWNER_JOB_IDS`` gets a pod-invariant offset, so +#: its replicas share one instant and the election settles it regardless of this bound. That +#: matters because the bound alone cannot be sized correctly for a job whose lease is released +#: when its body returns, which dedupes only for that body's runtime rather than for the lock's +#: TTL, so no non-zero spread would be safe. DEFAULT_CRON_DEDUPE_SECONDS: Final = MappingProxyType( { MONTHLY_SPEND_REPORT_JOB_ID: 3600, @@ -238,7 +247,7 @@ def _offset_for( return 0 return offset_seconds( job_id=job_id, - identity=identity, + identity="" if job_id in SINGLE_OWNER_JOB_IDS else identity, window_seconds=_window_for(job_id=job_id, period_seconds=period_seconds, settings=settings), ) diff --git a/litellm/proxy/common_utils/single_owner_job.py b/litellm/proxy/common_utils/single_owner_job.py index f04d372b815..2760379520b 100644 --- a/litellm/proxy/common_utils/single_owner_job.py +++ b/litellm/proxy/common_utils/single_owner_job.py @@ -13,11 +13,44 @@ from enum import Enum from typing import Final, TypeVar from litellm._logging import verbose_proxy_logger -from litellm.constants import SINGLE_OWNER_JOB_RENEWAL_DIVISOR +from litellm.constants import ( + EXPIRED_UI_SESSION_KEY_CLEANUP_JOB_NAME, + MONTHLY_SPEND_REPORT_JOB_ID, + PROMETHEUS_FALLBACK_STATS_JOB_ID, + PTU_ROLLUP_JOB_ID, + SINGLE_OWNER_JOB_RENEWAL_DIVISOR, + WEEKLY_SPEND_REPORT_JOB_ID, +) from litellm.proxy.db.db_transaction_queue.pod_lock_manager import PodLockManager T = TypeVar("T") +#: Scheduler ids of jobs that elect an owner before touching the database, so on any +#: given tick one replica does the work and the rest cost a single Redis read. +#: +#: These are the ids whose phase offset must NOT vary by pod. A per-pod offset spreads +#: replicas across the stagger window, and a lease released when its body returns only +#: dedupes for that body's runtime, so replicas further apart than that each find the +#: lock free and each run. Offsetting by job id alone keeps different jobs on different +#: instants, which is what staggering is for, while leaving every replica of one job on +#: the same instant, which is what lets the election do its job. +#: +#: Membership is "this job elects an owner", not "a serving pod skips it". The two differ: +#: the batch and responses cost pollers and the budget reset sweep are all role-gated while +#: taking no lock, so firing them together would have every replica do the work at once +#: rather than one. Add a job here only once its own entry point elects. +SINGLE_OWNER_JOB_IDS: Final = frozenset( + { + "spend_log_cleanup_job", + "key_rotation_job", + EXPIRED_UI_SESSION_KEY_CLEANUP_JOB_NAME, + PTU_ROLLUP_JOB_ID, + WEEKLY_SPEND_REPORT_JOB_ID, + MONTHLY_SPEND_REPORT_JOB_ID, + PROMETHEUS_FALLBACK_STATS_JOB_ID, + } +) + class JobRole(Enum): """Which scheduled jobs a process registers.""" diff --git a/tests/test_litellm/proxy/common_utils/test_scheduled_job_stagger.py b/tests/test_litellm/proxy/common_utils/test_scheduled_job_stagger.py index ca4d62737b6..6fc79336ee0 100644 --- a/tests/test_litellm/proxy/common_utils/test_scheduled_job_stagger.py +++ b/tests/test_litellm/proxy/common_utils/test_scheduled_job_stagger.py @@ -14,6 +14,7 @@ from apscheduler.triggers.interval import IntervalTrigger from litellm.constants import PTU_ROLLUP_JOB_ID, PTU_ROLLUP_LOCK_TTL_SECONDS from litellm.proxy._types import ScheduledJobStaggerSettings from litellm.proxy.common_utils.scheduled_job_stagger import ( + _offset_for, apply_scheduled_job_stagger, attach_job_timing_logger, offset_seconds, @@ -21,6 +22,7 @@ from litellm.proxy.common_utils.scheduled_job_stagger import ( resolve_stagger_identity, stagger_trigger, ) +from litellm.proxy.common_utils.single_owner_job import SINGLE_OWNER_JOB_IDS OPERATOR_CRON_JOB_ID = "spend_log_cleanup_job" SHARED_INTERVAL_JOB_IDS = ("periodic_reload_job", "get_credentials_job", "add_deployment_job") @@ -303,3 +305,78 @@ def test_job_timing_is_logged_with_scheduled_and_actual_start(caplog): assert f"scheduled_run_time={scheduled.isoformat()}" in message assert "actual_start_time=" in message assert "delay=2." in message + + +# --------------------------------------------------------------------------- +# Single-owner jobs: same instant on every replica, different instant per job + + +def _single_owner_offsets(job_id: str, **overrides) -> list[int]: + return [ + _offset_for( + job_id=job_id, + period_seconds=86400, + staggerable=True, + settings=_settings(**overrides), + identity=identity, + ) + for identity in ("pod-a:1", "pod-b:1", "pod-c:1", "pod-a:2") + ] + + +@pytest.mark.parametrize("job_id", sorted(SINGLE_OWNER_JOB_IDS)) +def test_every_replica_of_a_single_owner_job_shares_one_instant(job_id: str): + """An elected job must not be spread across replicas. + + Its lease is released when the body returns, so it dedupes for that body's runtime + and not for the lock's TTL. Replicas placed further apart than that each find the + lock free and each run, which is how a once-a-day job becomes once per replica per + day. Firing them together is what lets the election settle it. + """ + assert len(set(_single_owner_offsets(job_id))) == 1 + + +@pytest.mark.parametrize("job_id", ("update_spend_job", "periodic_reload_job", "update_gateway_requests_job")) +def test_a_per_pod_job_is_still_spread_across_replicas(job_id: str): + """Jobs that drain their own queues do the work on every pod, so they still need + separating; this is the property the single-owner case deliberately gives up.""" + assert len(set(_single_owner_offsets(job_id))) > 1 + + +def test_single_owner_jobs_still_land_on_different_instants_from_each_other(): + """Pod invariance must not collapse the jobs onto one instant: the elected replica + would then run every one of them at once, which is the burst being staggered away.""" + per_job = {job_id: _single_owner_offsets(job_id)[0] for job_id in SINGLE_OWNER_JOB_IDS} + + assert len(set(per_job.values())) == len(per_job) + + +@pytest.mark.parametrize( + "job_id", + ("check_batch_cost_job", "check_responses_cost_job", "reset_budget_job"), +) +def test_a_lockless_job_is_never_treated_as_single_owner(job_id: str): + """These are skipped by a serving pod and still take no lock. + + Membership is "elects an owner", not "a serving pod skips it". Pinning a lockless job + to one instant is strictly worse than spreading it, because every replica then does the + whole job at once instead of at staggered times. reset_budget_job is the sharpest case: + its sweep rewrites the entire due population, and synchronising that across a fleet is + the thundering herd the reset work exists to remove. It joins this set when its own + entry point elects, not before. + """ + assert job_id not in SINGLE_OWNER_JOB_IDS + assert len(set(_single_owner_offsets(job_id))) > 1, "a lockless job must stay spread across replicas" + + +def test_an_operator_offset_override_still_wins_for_a_single_owner_job(): + """The override is the operator's explicit instruction and predates this rule.""" + overridden = _offset_for( + job_id="key_rotation_job", + period_seconds=86400, + staggerable=True, + settings=_settings(offsets={"key_rotation_job": 17}), + identity="pod-a:1", + ) + + assert overridden == 17