- Add Model Mixins: TimestampMixin, SoftDeleteMixin, AuditMixin, etc. - Add Base Services: BaseService, BulkOperationsMixin, QueryOptimizerMixin - Add Base ViewSets with bulk operations - Add BackgroundJob model for Celery task tracking - Add BaseAppCommand for management commands - Add permissions, pagination, filters, cache, logging - Migrate tests to factory_boy + faker - Add CHANGELOG.md - 297 tests passing
238 lines
7.7 KiB
Python
238 lines
7.7 KiB
Python
"""
|
||
Background Job Tracking - отслеживание статуса Celery задач.
|
||
|
||
Предоставляет:
|
||
- Модель для хранения информации о задачах
|
||
- Сервис для управления задачами
|
||
- API эндпоинты для получения статуса
|
||
"""
|
||
|
||
import uuid
|
||
from typing import Any
|
||
|
||
from apps.core.mixins import TimestampMixin
|
||
from django.db import models
|
||
from django.utils import timezone
|
||
from django.utils.translation import gettext_lazy as _
|
||
|
||
|
||
class JobStatus(models.TextChoices):
|
||
"""Статусы фоновых задач."""
|
||
|
||
PENDING = "pending", _("Ожидает")
|
||
STARTED = "started", _("Выполняется")
|
||
SUCCESS = "success", _("Успешно")
|
||
FAILURE = "failure", _("Ошибка")
|
||
REVOKED = "revoked", _("Отменена")
|
||
RETRY = "retry", _("Повтор")
|
||
|
||
|
||
class BackgroundJob(TimestampMixin, models.Model):
|
||
"""
|
||
Модель для отслеживания фоновых задач Celery.
|
||
|
||
Позволяет:
|
||
- Отслеживать статус выполнения
|
||
- Хранить результат или ошибку
|
||
- Отображать прогресс выполнения
|
||
- Связывать задачу с пользователем
|
||
|
||
Использование в таске:
|
||
@shared_task(bind=True, base=TrackedTask)
|
||
def my_task(self, data):
|
||
job = BackgroundJob.objects.get(task_id=self.request.id)
|
||
job.update_progress(50, "Обработка...")
|
||
# ... логика
|
||
job.complete(result={"count": 100})
|
||
"""
|
||
|
||
id = models.UUIDField(
|
||
primary_key=True,
|
||
default=uuid.uuid4,
|
||
editable=False,
|
||
)
|
||
task_id = models.CharField(
|
||
_("ID задачи Celery"),
|
||
max_length=255,
|
||
unique=True,
|
||
db_index=True,
|
||
help_text=_("Идентификатор задачи в Celery"),
|
||
)
|
||
task_name = models.CharField(
|
||
_("имя задачи"),
|
||
max_length=255,
|
||
db_index=True,
|
||
help_text=_("Полное имя задачи (например, apps.myapp.tasks.process_data)"),
|
||
)
|
||
status = models.CharField(
|
||
_("статус"),
|
||
max_length=20,
|
||
choices=JobStatus.choices,
|
||
default=JobStatus.PENDING,
|
||
db_index=True,
|
||
)
|
||
progress = models.PositiveSmallIntegerField(
|
||
_("прогресс"),
|
||
default=0,
|
||
help_text=_("Прогресс выполнения в процентах (0-100)"),
|
||
)
|
||
progress_message = models.CharField(
|
||
_("сообщение о прогрессе"),
|
||
max_length=500,
|
||
blank=True,
|
||
default="",
|
||
)
|
||
result = models.JSONField(
|
||
_("результат"),
|
||
null=True,
|
||
blank=True,
|
||
help_text=_("Результат выполнения задачи (JSON)"),
|
||
)
|
||
error = models.TextField(
|
||
_("ошибка"),
|
||
blank=True,
|
||
default="",
|
||
help_text=_("Текст ошибки при неудачном выполнении"),
|
||
)
|
||
traceback = models.TextField(
|
||
_("traceback"),
|
||
blank=True,
|
||
default="",
|
||
help_text=_("Полный traceback ошибки"),
|
||
)
|
||
started_at = models.DateTimeField(
|
||
_("время начала"),
|
||
null=True,
|
||
blank=True,
|
||
)
|
||
completed_at = models.DateTimeField(
|
||
_("время завершения"),
|
||
null=True,
|
||
blank=True,
|
||
)
|
||
# Опционально: связь с пользователем
|
||
user_id = models.PositiveIntegerField(
|
||
_("ID пользователя"),
|
||
null=True,
|
||
blank=True,
|
||
db_index=True,
|
||
help_text=_("ID пользователя, запустившего задачу"),
|
||
)
|
||
# Метаданные
|
||
meta = models.JSONField(
|
||
_("метаданные"),
|
||
default=dict,
|
||
blank=True,
|
||
help_text=_("Дополнительные данные задачи"),
|
||
)
|
||
|
||
class Meta:
|
||
verbose_name = _("фоновая задача")
|
||
verbose_name_plural = _("фоновые задачи")
|
||
ordering = ["-created_at"]
|
||
indexes = [
|
||
models.Index(fields=["status", "created_at"]),
|
||
models.Index(fields=["user_id", "status"]),
|
||
models.Index(fields=["task_name", "status"]),
|
||
]
|
||
|
||
def __str__(self) -> str:
|
||
return f"{self.task_name} ({self.status})"
|
||
|
||
# ==================== Методы обновления статуса ====================
|
||
|
||
def mark_started(self) -> None:
|
||
"""Отметить задачу как начатую."""
|
||
self.status = JobStatus.STARTED
|
||
self.started_at = timezone.now()
|
||
self.save(update_fields=["status", "started_at", "updated_at"])
|
||
|
||
def update_progress(self, progress: int, message: str = "") -> None:
|
||
"""
|
||
Обновить прогресс выполнения.
|
||
|
||
Args:
|
||
progress: Процент выполнения (0-100)
|
||
message: Описание текущего этапа
|
||
"""
|
||
self.progress = min(max(progress, 0), 100)
|
||
self.progress_message = message
|
||
self.save(update_fields=["progress", "progress_message", "updated_at"])
|
||
|
||
def complete(self, result: Any = None) -> None:
|
||
"""
|
||
Отметить задачу как успешно завершённую.
|
||
|
||
Args:
|
||
result: Результат выполнения (сериализуемый в JSON)
|
||
"""
|
||
self.status = JobStatus.SUCCESS
|
||
self.progress = 100
|
||
self.result = result
|
||
self.completed_at = timezone.now()
|
||
self.save(
|
||
update_fields=[
|
||
"status",
|
||
"progress",
|
||
"result",
|
||
"completed_at",
|
||
"updated_at",
|
||
]
|
||
)
|
||
|
||
def fail(self, error: str, traceback_str: str = "") -> None:
|
||
"""
|
||
Отметить задачу как завершённую с ошибкой.
|
||
|
||
Args:
|
||
error: Текст ошибки
|
||
traceback_str: Полный traceback
|
||
"""
|
||
self.status = JobStatus.FAILURE
|
||
self.error = str(error)
|
||
self.traceback = traceback_str
|
||
self.completed_at = timezone.now()
|
||
self.save(
|
||
update_fields=[
|
||
"status",
|
||
"error",
|
||
"traceback",
|
||
"completed_at",
|
||
"updated_at",
|
||
]
|
||
)
|
||
|
||
def revoke(self) -> None:
|
||
"""Отметить задачу как отменённую."""
|
||
self.status = JobStatus.REVOKED
|
||
self.completed_at = timezone.now()
|
||
self.save(update_fields=["status", "completed_at", "updated_at"])
|
||
|
||
def mark_retry(self) -> None:
|
||
"""Отметить, что задача будет повторена."""
|
||
self.status = JobStatus.RETRY
|
||
self.save(update_fields=["status", "updated_at"])
|
||
|
||
# ==================== Свойства ====================
|
||
|
||
@property
|
||
def is_finished(self) -> bool:
|
||
"""Проверка, завершена ли задача."""
|
||
return self.status in (
|
||
JobStatus.SUCCESS,
|
||
JobStatus.FAILURE,
|
||
JobStatus.REVOKED,
|
||
)
|
||
|
||
@property
|
||
def is_successful(self) -> bool:
|
||
"""Проверка успешного завершения."""
|
||
return self.status == JobStatus.SUCCESS
|
||
|
||
@property
|
||
def duration(self) -> float | None:
|
||
"""Длительность выполнения в секундах."""
|
||
if self.started_at and self.completed_at:
|
||
return (self.completed_at - self.started_at).total_seconds()
|
||
return None
|