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

1# Copyright © Michal Čihař <michal@weblate.org> 

2# 

3# SPDX-License-Identifier: GPL-3.0-or-later 

4 

5"""Celery integration helper tools.""" 

6 

7from __future__ import annotations 

8 

9import os 

10import time 

11from collections import defaultdict 

12 

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 

18 

19# set the default Django settings module for the 'celery' program. 

20os.environ.setdefault("DJANGO_SETTINGS_MODULE", "weblate.settings") 

21 

22app = Celery("weblate") 

23 

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") 

29 

30# Load task modules from all registered Django app configs. 

31app.autodiscover_tasks() 

32 

33 

34@task_failure.connect 

35def handle_task_failure(task_id="", exception=None, **kwargs) -> None: 

36 from weblate.utils.errors import report_error 

37 

38 report_error( 

39 f"Failure while executing task {task_id}", 

40 skip_sentry=True, 

41 print_tb=True, 

42 level="error", 

43 ) 

44 

45 

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 

50 

51 init_error_collection(celery=True) 

52 

53 

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) 

63 

64 

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 

70 

71 

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 

79 

80 

81def get_queue_stats(): 

82 """Calculate queue stats.""" 

83 return {queue: get_queue_length(queue) for queue in get_queue_list()} 

84 

85 

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"] 

95 

96 # Not yet started 

97 return 0 

98 

99 

100def is_celery_queue_long(): 

101 """ 

102 Check whether celery queue is too long. 

103 

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 

109 

110 cache_key = "celery_queue_stats" 

111 queues_data = cache.get(cache_key, {}) 

112 

113 # Hours since epoch 

114 current_hour = int(time.time() / 3600) 

115 test_hour = current_hour - 1 

116 

117 # Fetch current stats 

118 stats = get_queue_stats() 

119 

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 

128 

129 # Store to cache 

130 cache.set(cache_key, queues_data, 7200) 

131 

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 

135 

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 )