v0.163.0
  1from __future__ import annotations
  2
  3from datetime import timedelta
  4from typing import ClassVar
  5
  6from plain.admin.cards import Card, TrendCard
  7from plain.admin.views import (
  8    AdminModelDetailView,
  9    AdminModelListView,
 10    AdminViewset,
 11    register_viewset,
 12)
 13from plain.http import RedirectResponse
 14from plain.postgres.expressions import Case, When
 15from plain.runtime import settings
 16
 17from plain import postgres
 18
 19from .models import (
 20    JobProcess,
 21    JobRequest,
 22    JobResult,
 23    JobResultQuerySet,
 24    WorkerHeartbeat,
 25    heartbeat_cutoff,
 26)
 27
 28
 29def _td_format(td_object: timedelta) -> str:
 30    seconds = int(td_object.total_seconds())
 31    periods = [
 32        ("year", 60 * 60 * 24 * 365),
 33        ("month", 60 * 60 * 24 * 30),
 34        ("day", 60 * 60 * 24),
 35        ("hour", 60 * 60),
 36        ("minute", 60),
 37        ("second", 1),
 38    ]
 39
 40    strings = []
 41    for period_name, period_seconds in periods:
 42        if seconds > period_seconds:
 43            period_value, seconds = divmod(seconds, period_seconds)
 44            has_s = "s" if period_value > 1 else ""
 45            strings.append(f"{period_value} {period_name}{has_s}")
 46
 47    return ", ".join(strings)
 48
 49
 50class JobResultsTrendCard(TrendCard):
 51    title = "Results trend"
 52    model = JobResult
 53    datetime_field = "created_at"
 54    size = TrendCard.Sizes.FULL
 55    group_field = "status"
 56    group_labels: ClassVar = {
 57        "SUCCESSFUL": "Successful",
 58        "ERRORED": "Errored",
 59        "CANCELLED": "Cancelled",
 60        "DEFERRED": "Deferred",
 61        "LOST": "Lost",
 62    }
 63    group_colors: ClassVar = {
 64        "SUCCESSFUL": "var(--success)",
 65        "ERRORED": "var(--danger)",
 66        "CANCELLED": "var(--muted-foreground)",
 67        "DEFERRED": "var(--info)",
 68        "LOST": "var(--warning)",
 69    }
 70
 71
 72class SuccessfulJobsCard(Card):
 73    title = "Successful"
 74    text = "View"
 75
 76    def get_metric(self) -> int:
 77        return JobResult.query.successful().count()
 78
 79    def get_link(self) -> str:
 80        return JobResultViewset.ListView.get_view_url() + "?display=Successful"
 81
 82
 83class ErroredJobsCard(Card):
 84    title = "Errored"
 85    text = "View"
 86
 87    def get_metric(self) -> int:
 88        return JobResult.query.errored().count()
 89
 90    def get_link(self) -> str:
 91        return JobResultViewset.ListView.get_view_url() + "?display=Errored"
 92
 93
 94class LostJobsCard(Card):
 95    title = "Lost"
 96    text = "View"  # TODO make not required - just an icon?
 97
 98    def get_description(self) -> str:
 99        delta = timedelta(seconds=settings.JOBS_HEARTBEAT_TIMEOUT)
100        return (
101            f"Jobs are considered lost when their worker stops heartbeating "
102            f"for more than {_td_format(delta)}"
103        )
104
105    def get_metric(self) -> int:
106        return JobResult.query.lost().count()
107
108    def get_link(self) -> str:
109        return JobResultViewset.ListView.get_view_url() + "?display=Lost"
110
111
112class RetriedJobsCard(Card):
113    title = "Retried"
114    text = "View"  # TODO make not required - just an icon?
115
116    def get_metric(self) -> int:
117        return JobResult.query.retried().count()
118
119    def get_link(self) -> str:
120        return JobResultViewset.ListView.get_view_url() + "?display=Retried"
121
122
123class WaitingJobsCard(Card):
124    title = "Waiting"
125
126    def get_metric(self) -> int:
127        return JobProcess.query.waiting().count()
128
129
130class RunningJobsCard(Card):
131    title = "Running"
132
133    def get_metric(self) -> int:
134        return JobProcess.query.running().count()
135
136
137class ActiveWorkersCard(Card):
138    title = "Active workers"
139    text = "View"
140
141    def get_description(self) -> str:
142        delta = timedelta(seconds=settings.JOBS_HEARTBEAT_TIMEOUT)
143        return f"Workers whose heartbeat is within the last {_td_format(delta)}."
144
145    def get_metric(self) -> int:
146        return WorkerHeartbeat.query.filter(
147            last_heartbeat_at__gte=heartbeat_cutoff()
148        ).count()
149
150    def get_link(self) -> str:
151        return WorkerHeartbeatViewset.ListView.get_view_url()
152
153
154class StaleWorkersCard(Card):
155    title = "Stale workers"
156    text = "View"
157
158    def get_description(self) -> str:
159        delta = timedelta(seconds=settings.JOBS_HEARTBEAT_TIMEOUT)
160        return (
161            f"Workers whose heartbeat is older than {_td_format(delta)}. "
162            f"Their in-flight jobs are about to be rescued as Lost."
163        )
164
165    def get_metric(self) -> int:
166        return WorkerHeartbeat.query.filter(
167            last_heartbeat_at__lt=heartbeat_cutoff()
168        ).count()
169
170    def get_link(self) -> str:
171        return WorkerHeartbeatViewset.ListView.get_view_url() + "?display=Stale"
172
173
174@register_viewset
175class JobRequestViewset(AdminViewset):
176    class ListView(AdminModelListView):
177        nav_section = "Jobs"
178        nav_icon = "inbox"
179        model = JobRequest
180        title = "Requests"
181        description = "Jobs waiting to be picked up by a worker."
182        fields = (
183            "id",
184            "job_class",
185            "priority",
186            "created_at",
187            "start_at",
188            "concurrency_key",
189        )
190        actions = ("Delete",)
191        queryset_order = ("-priority", "-start_at", "-created_at")
192
193        def perform_action(self, action: str, objects: postgres.QuerySet) -> None:
194            if action == "Delete":
195                objects.delete()
196
197    class DetailView(AdminModelDetailView):
198        model = JobRequest
199        title = "Request"
200
201
202@register_viewset
203class JobProcessViewset(AdminViewset):
204    class ListView(AdminModelListView):
205        nav_section = "Jobs"
206        nav_icon = "gear"
207        model = JobProcess
208        title = "Processes"
209        description = "Jobs currently being processed by a worker."
210        fields = (
211            "id",
212            "job_class",
213            "priority",
214            "created_at",
215            "started_at",
216            "concurrency_key",
217        )
218        actions = ("Delete",)
219        cards = (
220            WaitingJobsCard,
221            RunningJobsCard,
222            ActiveWorkersCard,
223            StaleWorkersCard,
224        )
225
226        def perform_action(self, action: str, objects: postgres.QuerySet) -> None:
227            if action == "Delete":
228                objects.delete()
229
230    class DetailView(AdminModelDetailView):
231        model = JobProcess
232        title = "Process"
233
234
235@register_viewset
236class JobResultViewset(AdminViewset):
237    class ListView(AdminModelListView):
238        nav_section = "Jobs"
239        nav_icon = "clipboard-check"
240        model = JobResult
241        title = "Results"
242        description = "Completed jobs with their success/failure status."
243        fields = (
244            "id",
245            "job_class",
246            "priority",
247            "created_at",
248            "status",
249            "retried",
250            "is_retry",
251        )
252        field_templates: ClassVar = {
253            "status": "jobs/values/job_status.html",
254        }
255        search_fields = (
256            "uuid",
257            "job_process_uuid",
258            "job_request_uuid",
259            "job_class",
260        )
261        cards = (
262            JobResultsTrendCard,
263            SuccessfulJobsCard,
264            ErroredJobsCard,
265            LostJobsCard,
266            RetriedJobsCard,
267        )
268        filters = (
269            "Successful",
270            "Errored",
271            "Cancelled",
272            "Lost",
273            "Retried",
274        )
275        actions = ("Retry",)
276
277        def get_initial_queryset(self) -> JobResultQuerySet:
278            queryset: JobResultQuerySet = super().get_initial_queryset()  # ty: ignore[invalid-assignment]
279            return queryset.annotate(
280                retried=Case(
281                    When(retry_job_request_uuid__isnull=False, then=True),
282                    default=False,
283                    output_field=postgres.BooleanField(),
284                ),
285                is_retry=Case(
286                    When(retry_attempt__gt=0, then=True),
287                    default=False,
288                    output_field=postgres.BooleanField(),
289                ),
290            )
291
292        def filter_queryset(self, queryset: JobResultQuerySet) -> JobResultQuerySet:
293            if self.filter == "Successful":
294                return queryset.successful()
295            if self.filter == "Errored":
296                return queryset.errored()
297            if self.filter == "Cancelled":
298                return queryset.cancelled()
299            if self.filter == "Lost":
300                return queryset.lost()
301            if self.filter == "Retried":
302                return queryset.retried()
303            return queryset
304
305        def get_fields(self) -> tuple[str, ...]:
306            fields = super().get_fields()
307            if self.filter == "Retried":
308                fields = (*fields, "retries", "retry_attempt")
309            return fields
310
311        def perform_action(self, action: str, objects: postgres.QuerySet) -> None:
312            if action == "Retry":
313                for result in objects:
314                    result.retry_job(delay=0)
315            else:
316                raise ValueError("Invalid action")
317
318    class DetailView(AdminModelDetailView):
319        model = JobResult
320        title = "Result"
321
322        def post(self) -> RedirectResponse:
323            self.object.retry_job(delay=0)
324            return RedirectResponse(self.request.get_full_path(), status_code=302)
325
326
327@register_viewset
328class WorkerHeartbeatViewset(AdminViewset):
329    class ListView(AdminModelListView):
330        nav_section = "Jobs"
331        nav_icon = "heart-pulse"
332        model = WorkerHeartbeat
333        title = "Workers"
334        description = (
335            "Live worker processes. Each row is refreshed while its worker is "
336            "running and deleted on clean shutdown."
337        )
338        fields = (
339            "worker_id",
340            "hostname",
341            "pid",
342            "queues",
343            "started_at",
344            "last_heartbeat_at",
345            "stale",
346        )
347        search_fields = (
348            "worker_id",
349            "hostname",
350        )
351        filters = (
352            "Active",
353            "Stale",
354        )
355        queryset_order = ("-last_heartbeat_at",)
356
357        def get_initial_queryset(self) -> postgres.QuerySet[WorkerHeartbeat]:
358            queryset = super().get_initial_queryset()
359            return queryset.annotate(
360                stale=Case(
361                    When(last_heartbeat_at__lt=heartbeat_cutoff(), then=True),
362                    default=False,
363                    output_field=postgres.BooleanField(),
364                ),
365            )
366
367        def filter_queryset(
368            self, queryset: postgres.QuerySet[WorkerHeartbeat]
369        ) -> postgres.QuerySet[WorkerHeartbeat]:
370            cutoff = heartbeat_cutoff()
371            if self.filter == "Active":
372                return queryset.filter(last_heartbeat_at__gte=cutoff)
373            if self.filter == "Stale":
374                return queryset.filter(last_heartbeat_at__lt=cutoff)
375            return queryset
376
377    class DetailView(AdminModelDetailView):
378        model = WorkerHeartbeat
379        title = "Worker"