112 lines
4.0 KiB
Python
112 lines
4.0 KiB
Python
from django.db import connection
|
|
from rest_framework import generics, status
|
|
from rest_framework.exceptions import NotFound, ValidationError
|
|
from rest_framework.pagination import CursorPagination, LimitOffsetPagination
|
|
from rest_framework.response import Response
|
|
from rest_framework.views import APIView
|
|
|
|
from jobs.models import Job, JobEvent
|
|
from jobs.serializers import JobCreateSerializer, JobEventSerializer, JobSerializer
|
|
from jobs.services import create_job, job_stats, retry_failed_job
|
|
|
|
|
|
class JobEventCursorPagination(CursorPagination):
|
|
page_size = 50
|
|
page_size_query_param = "limit"
|
|
max_page_size = 200
|
|
ordering = "-id"
|
|
|
|
|
|
class JobLimitOffsetPagination(LimitOffsetPagination):
|
|
default_limit = 25
|
|
max_limit = 100
|
|
|
|
|
|
class HealthAPIView(APIView):
|
|
def get(self, request):
|
|
try:
|
|
with connection.cursor() as cursor:
|
|
cursor.execute("SELECT 1")
|
|
cursor.fetchone()
|
|
except Exception as exc:
|
|
return Response(
|
|
{"ok": False, "database": "error", "detail": str(exc)},
|
|
status=status.HTTP_503_SERVICE_UNAVAILABLE,
|
|
)
|
|
return Response({"ok": True, "database": "ok"})
|
|
|
|
|
|
class JobListCreateAPIView(APIView):
|
|
def get(self, request):
|
|
queryset = Job.objects.order_by("-created_at")
|
|
status_filter = request.query_params.get("status")
|
|
type_filter = request.query_params.get("type")
|
|
if status_filter:
|
|
queryset = queryset.filter(status=status_filter)
|
|
if type_filter:
|
|
queryset = queryset.filter(type__icontains=type_filter)
|
|
|
|
paginator = JobLimitOffsetPagination()
|
|
page = paginator.paginate_queryset(queryset, request, view=self)
|
|
return paginator.get_paginated_response(JobSerializer(page, many=True).data)
|
|
|
|
def post(self, request):
|
|
serializer = JobCreateSerializer(data=request.data)
|
|
serializer.is_valid(raise_exception=True)
|
|
result = create_job(
|
|
job_type=serializer.validated_data["type"],
|
|
payload=serializer.validated_data.get("payload") or {},
|
|
priority=serializer.validated_data.get("priority", 50),
|
|
available_at=serializer.validated_data.get("available_at"),
|
|
max_attempts=serializer.validated_data.get("max_attempts"),
|
|
idempotency_key=serializer.validated_data.get("idempotency_key"),
|
|
)
|
|
return Response(
|
|
JobSerializer(result.job).data,
|
|
status=status.HTTP_201_CREATED if result.created else status.HTTP_200_OK,
|
|
)
|
|
|
|
|
|
class JobDetailAPIView(generics.RetrieveAPIView):
|
|
queryset = Job.objects.all()
|
|
serializer_class = JobSerializer
|
|
|
|
|
|
class JobEventsAPIView(APIView):
|
|
def get(self, request, pk):
|
|
if not Job.objects.filter(id=pk).exists():
|
|
raise NotFound("Job not found.")
|
|
events = JobEvent.objects.filter(job_id=pk).order_by("id")
|
|
return Response(JobEventSerializer(events, many=True).data)
|
|
|
|
|
|
class GlobalJobEventsAPIView(APIView):
|
|
def get(self, request):
|
|
queryset = JobEvent.objects.order_by("-id")
|
|
event_type = request.query_params.get("type")
|
|
job_id = request.query_params.get("job_id")
|
|
if event_type:
|
|
queryset = queryset.filter(type=event_type)
|
|
if job_id:
|
|
queryset = queryset.filter(job_id=job_id)
|
|
|
|
paginator = JobEventCursorPagination()
|
|
page = paginator.paginate_queryset(queryset, request, view=self)
|
|
return paginator.get_paginated_response(JobEventSerializer(page, many=True).data)
|
|
|
|
|
|
class JobStatsAPIView(APIView):
|
|
def get(self, request):
|
|
return Response(job_stats())
|
|
|
|
|
|
class RetryJobAPIView(APIView):
|
|
def post(self, request, pk):
|
|
try:
|
|
job = retry_failed_job(pk)
|
|
except Job.DoesNotExist as exc:
|
|
raise NotFound("Job not found.") from exc
|
|
except ValueError as exc:
|
|
raise ValidationError({"detail": str(exc)}) from exc
|
|
return Response(JobSerializer(job).data)
|