Coverage for app/venv/lib/python3.14/site-packages/weblate/utils/celery.py: 79%
65 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 07:15 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 07:15 +0000
1# Copyright © Michal Čihař <michal@weblate.org>
2#
3# SPDX-License-Identifier: GPL-3.0-or-later
5"""Celery integration helper tools."""
7from __future__ import annotations
9import os
10import time
11from collections import defaultdict
13from celery import Celery
14from celery.signals import after_setup_logger, task_failure
15from django.conf import settings
16from django.core.cache import cache
17from django.core.checks import run_checks
19# set the default Django settings module for the 'celery' program.
20os.environ.setdefault("DJANGO_SETTINGS_MODULE", "weblate.settings")
22app = Celery("weblate")
24# Using a string here means the worker doesn't have to serialize
25# the configuration object to child processes.
26# - namespace='CELERY' means all celery-related configuration keys
27# should have a `CELERY_` prefix.
28app.config_from_object("django.conf:settings", namespace="CELERY")
30# Load task modules from all registered Django app configs.
31app.autodiscover_tasks()
34@task_failure.connect
35def handle_task_failure(task_id="", exception=None, **kwargs) -> None:
36 from weblate.utils.errors import report_error
38 report_error(
39 f"Failure while executing task {task_id}",
40 skip_sentry=True,
41 print_tb=True,
42 level="error",
43 )
46@app.on_after_configure.connect
47def configure_error_handling(sender, **kargs) -> None:
48 """Rollbar and Sentry integration."""
49 from weblate.utils.errors import init_error_collection
51 init_error_collection(celery=True)
54@after_setup_logger.connect
55def show_failing_system_check(sender, logger, **kwargs) -> None:
56 if settings.DEBUG: 56 ↛ 57line 56 didn't jump to line 57 because the condition on line 56 was never true
57 for check in run_checks(include_deployment_checks=True):
58 # Skip silenced checks and Celery one
59 # (it fails when started from Celery startup)
60 if check.is_silenced() or check.id == "weblate.E019":
61 continue
62 logger.warning("%s", check)
65def get_queue_length(queue="celery"):
66 with app.connection_or_acquire() as conn: # type: ignore[attr-defined]
67 return conn.default_channel.queue_declare(
68 queue=queue, durable=True, auto_delete=False
69 ).message_count
72def get_queue_list():
73 """List queues in Celery."""
74 result = {"celery"}
75 for route in settings.CELERY_TASK_ROUTES.values():
76 if "queue" in route: 76 ↛ 75line 76 didn't jump to line 75 because the condition on line 76 was always true
77 result.add(route["queue"])
78 return result
81def get_queue_stats():
82 """Calculate queue stats."""
83 return {queue: get_queue_length(queue) for queue in get_queue_list()}
86def get_task_progress(task):
87 """Return progress of a Celery task."""
88 # Completed task
89 if task.ready(): 89 ↛ 90line 89 didn't jump to line 90 because the condition on line 89 was never true
90 return 100
91 # In progress
92 result = task.result
93 if task.state == "PROGRESS" and result is not None: 93 ↛ 94line 93 didn't jump to line 94 because the condition on line 93 was never true
94 return result["progress"]
96 # Not yet started
97 return 0
100def is_celery_queue_long():
101 """
102 Check whether celery queue is too long.
104 It does trigger if it is too long for at least one hour. This way peaks are
105 filtered out, and no warning need be issued for big operations (for example
106 site-wide autotranslation).
107 """
108 from weblate.trans.models import Translation
110 cache_key = "celery_queue_stats"
111 queues_data = cache.get(cache_key, {})
113 # Hours since epoch
114 current_hour = int(time.time() / 3600)
115 test_hour = current_hour - 1
117 # Fetch current stats
118 stats = get_queue_stats()
120 # Update counters
121 if current_hour not in queues_data:
122 # Delete stale items
123 for key in list(queues_data.keys()):
124 if key < test_hour: 124 ↛ 125line 124 didn't jump to line 125 because the condition on line 124 was never true
125 del queues_data[key]
126 # Add current one
127 queues_data[current_hour] = stats
129 # Store to cache
130 cache.set(cache_key, queues_data, 7200)
132 # Do not fire if we do not have counts for two hours ago
133 if test_hour not in queues_data:
134 return False
136 # Check if any queue got bigger
137 base = queues_data[test_hour]
138 thresholds: dict[str, int] = defaultdict(lambda: 50)
139 # Set the limit to avoid trigger on auto-translating all components
140 # nightly.
141 thresholds["translate"] = max(1000, Translation.objects.count() // 30)
142 return any(
143 stat > thresholds[key] and base.get(key, 0) > thresholds[key]
144 for key, stat in stats.items()
145 )