2022-09-29 12:44:45 +00:00
|
|
|
import ast
|
|
|
|
|
2022-12-22 09:49:07 +00:00
|
|
|
from celery import signals
|
2022-11-21 11:54:00 +00:00
|
|
|
from django.core.cache import cache
|
2022-12-22 09:49:07 +00:00
|
|
|
from django.db import transaction
|
2022-12-20 11:46:19 +00:00
|
|
|
from django.db.models.signals import pre_save
|
2022-11-21 11:54:00 +00:00
|
|
|
from django.db.utils import ProgrammingError
|
2022-12-22 09:49:07 +00:00
|
|
|
from django.dispatch import receiver
|
2022-09-29 12:44:45 +00:00
|
|
|
from django.utils import translation, timezone
|
|
|
|
|
2022-11-21 11:54:00 +00:00
|
|
|
from common.db.utils import close_old_connections, get_logger
|
2023-02-17 09:14:53 +00:00
|
|
|
from common.signals import django_ready
|
2023-02-20 07:01:00 +00:00
|
|
|
from orgs.utils import get_current_org_id, set_current_org
|
2022-10-24 12:14:18 +00:00
|
|
|
from .celery import app
|
2022-12-20 11:46:19 +00:00
|
|
|
from .models import CeleryTaskExecution, CeleryTask, Job
|
2022-01-13 06:18:32 +00:00
|
|
|
|
2022-09-29 12:44:45 +00:00
|
|
|
logger = get_logger(__name__)
|
2022-01-13 06:18:32 +00:00
|
|
|
|
|
|
|
|
2022-12-20 11:46:19 +00:00
|
|
|
@receiver(pre_save, sender=Job)
|
|
|
|
def on_account_pre_create(sender, instance, **kwargs):
|
|
|
|
# 升级版本号
|
|
|
|
instance.version += 1
|
|
|
|
|
|
|
|
|
2022-12-22 09:49:07 +00:00
|
|
|
@receiver(signals.worker_ready)
|
2022-10-24 12:14:18 +00:00
|
|
|
def sync_registered_tasks(*args, **kwargs):
|
2022-12-22 09:49:07 +00:00
|
|
|
synced = cache.get('synced_registered_tasks', False)
|
|
|
|
if synced:
|
|
|
|
return
|
|
|
|
cache.set('synced_registered_tasks', True, 60)
|
2022-10-24 12:14:18 +00:00
|
|
|
with transaction.atomic():
|
2022-10-26 09:38:32 +00:00
|
|
|
try:
|
|
|
|
db_tasks = CeleryTask.objects.all()
|
2022-11-21 11:54:00 +00:00
|
|
|
celery_task_names = [key for key in app.tasks]
|
|
|
|
db_task_names = db_tasks.values_list('name', flat=True)
|
2022-10-26 09:38:32 +00:00
|
|
|
|
2022-11-21 11:54:00 +00:00
|
|
|
db_tasks.exclude(name__in=celery_task_names).delete()
|
|
|
|
not_in_db_tasks = set(celery_task_names) - set(db_task_names)
|
|
|
|
tasks_to_create = [CeleryTask(name=name) for name in not_in_db_tasks]
|
|
|
|
CeleryTask.objects.bulk_create(tasks_to_create)
|
|
|
|
except ProgrammingError:
|
|
|
|
pass
|
2022-10-24 12:14:18 +00:00
|
|
|
|
|
|
|
|
2023-02-17 09:14:53 +00:00
|
|
|
@receiver(django_ready)
|
|
|
|
def check_registered_tasks(*args, **kwargs):
|
|
|
|
attrs = ['verbose_name', 'activity_callback']
|
|
|
|
for name, task in app.tasks.items():
|
|
|
|
if name.startswith('celery.'):
|
|
|
|
continue
|
|
|
|
for attr in attrs:
|
|
|
|
if not hasattr(task, attr):
|
2023-02-20 11:12:57 +00:00
|
|
|
# print('>>> Task {} has no attribute {}'.format(name, attr))
|
|
|
|
pass
|
2023-02-17 09:14:53 +00:00
|
|
|
|
|
|
|
|
2022-09-29 12:44:45 +00:00
|
|
|
@signals.before_task_publish.connect
|
2023-02-20 07:01:00 +00:00
|
|
|
def before_task_publish(body=None, **kwargs):
|
2022-01-13 06:18:32 +00:00
|
|
|
current_lang = translation.get_language()
|
2023-02-20 07:01:00 +00:00
|
|
|
current_org_id = get_current_org_id()
|
|
|
|
args, kwargs = body[:2]
|
|
|
|
kwargs['__current_lang'] = current_lang
|
|
|
|
kwargs['__current_org_id'] = current_org_id
|
2022-01-13 06:18:32 +00:00
|
|
|
|
|
|
|
|
2022-09-29 12:44:45 +00:00
|
|
|
@signals.task_prerun.connect
|
2023-02-20 07:01:00 +00:00
|
|
|
def on_celery_task_pre_run(task_id='', kwargs=None, **others):
|
2022-09-29 12:44:45 +00:00
|
|
|
# 更新状态
|
2022-11-21 11:54:00 +00:00
|
|
|
CeleryTaskExecution.objects.filter(id=task_id) \
|
2022-11-01 09:04:44 +00:00
|
|
|
.update(state='RUNNING', date_start=timezone.now())
|
2022-01-17 03:07:27 +00:00
|
|
|
# 关闭之前的数据库连接
|
|
|
|
close_old_connections()
|
|
|
|
|
2023-02-20 07:01:00 +00:00
|
|
|
# 设置语言的一些上下文
|
|
|
|
lang = kwargs.pop('__current_lang', None)
|
|
|
|
org_id = kwargs.pop('__current_org_id', None)
|
|
|
|
if lang:
|
|
|
|
print('>> Set language to {}'.format(lang))
|
|
|
|
translation.activate(lang)
|
|
|
|
if org_id:
|
|
|
|
print('>> Set org to {}'.format(org_id))
|
|
|
|
set_current_org(org_id)
|
2022-01-17 03:07:27 +00:00
|
|
|
|
|
|
|
|
2022-09-29 12:44:45 +00:00
|
|
|
@signals.task_postrun.connect
|
|
|
|
def on_celery_task_post_run(task_id='', state='', **kwargs):
|
2022-01-17 03:07:27 +00:00
|
|
|
close_old_connections()
|
2022-09-29 12:44:45 +00:00
|
|
|
|
2022-10-24 12:14:18 +00:00
|
|
|
CeleryTaskExecution.objects.filter(id=task_id).update(
|
2022-09-29 12:44:45 +00:00
|
|
|
state=state, date_finished=timezone.now(), is_finished=True
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
@signals.after_task_publish.connect
|
|
|
|
def task_sent_handler(headers=None, body=None, **kwargs):
|
|
|
|
info = headers if 'task' in headers else body
|
|
|
|
task = info.get('task')
|
|
|
|
i = info.get('id')
|
|
|
|
if not i or not task:
|
|
|
|
logger.error("Not found task id or name: {}".format(info))
|
|
|
|
return
|
|
|
|
|
|
|
|
args = info.get('argsrepr', '()')
|
|
|
|
kwargs = info.get('kwargsrepr', '{}')
|
|
|
|
try:
|
|
|
|
args = list(ast.literal_eval(args))
|
|
|
|
kwargs = ast.literal_eval(kwargs)
|
|
|
|
except (ValueError, SyntaxError):
|
|
|
|
args = []
|
|
|
|
kwargs = {}
|
|
|
|
|
|
|
|
data = {
|
|
|
|
'id': i,
|
|
|
|
'name': task,
|
|
|
|
'state': 'PENDING',
|
|
|
|
'is_finished': False,
|
|
|
|
'args': args,
|
|
|
|
'kwargs': kwargs
|
|
|
|
}
|
2022-10-24 12:14:18 +00:00
|
|
|
CeleryTaskExecution.objects.create(**data)
|
2023-01-16 11:02:09 +00:00
|
|
|
CeleryTask.objects.filter(name=task).update(date_last_publish=timezone.now())
|