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"