1from __future__ import annotations
2
3import datetime
4import signal
5from typing import Any
6
7import click
8from plain.cli import SettingOption, register_cli
9from plain.logs import get_framework_logger
10from plain.runtime import settings
11from plain.utils import timezone
12
13from .models import JobProcess, JobRequest, JobResult
14from .registry import jobs_registry
15from .scheduling import load_schedule
16from .workers import Worker
17
18logger = get_framework_logger()
19
20
21@register_cli("jobs")
22@click.group()
23def cli() -> None:
24 """Background job management"""
25
26
27@cli.command()
28@click.option(
29 "queues",
30 "--queue",
31 default=["default"],
32 multiple=True,
33 type=str,
34 help="Queue to process",
35)
36@click.option(
37 "--max-processes",
38 "max_processes",
39 type=int,
40 cls=SettingOption,
41 setting="JOBS_WORKER_MAX_PROCESSES",
42)
43@click.option(
44 "--max-jobs-per-process",
45 "max_jobs_per_process",
46 type=int,
47 cls=SettingOption,
48 setting="JOBS_WORKER_MAX_JOBS_PER_PROCESS",
49)
50@click.option(
51 "--max-pending-per-process",
52 "max_pending_per_process",
53 type=int,
54 cls=SettingOption,
55 setting="JOBS_WORKER_MAX_PENDING_PER_PROCESS",
56)
57@click.option(
58 "--stats-every",
59 "stats_every",
60 type=int,
61 cls=SettingOption,
62 setting="JOBS_WORKER_STATS_EVERY",
63)
64@click.option(
65 "--reload",
66 is_flag=True,
67 help="Watch files and auto-reload worker on changes",
68)
69def worker(
70 queues: tuple[str, ...],
71 max_processes: int | None,
72 max_jobs_per_process: int | None,
73 max_pending_per_process: int,
74 stats_every: int,
75 reload: bool,
76) -> None:
77 """Run the job worker"""
78 jobs_schedule = load_schedule(settings.JOBS_SCHEDULE)
79
80 worker_kwargs = {
81 "queues": list(queues),
82 "jobs_schedule": jobs_schedule,
83 "max_processes": max_processes,
84 "max_jobs_per_process": max_jobs_per_process,
85 "max_pending_per_process": max_pending_per_process,
86 "stats_every": stats_every,
87 }
88
89 if reload:
90 _run_with_reload(worker_kwargs)
91 else:
92 _run_once(worker_kwargs)
93
94
95def _run_with_reload(worker_kwargs: dict[str, Any]) -> None:
96 from plain.internal.reloader import Reloader
97
98 should_restart = {"value": True}
99 current_worker: dict[str, Worker | None] = {"instance": None}
100
101 def file_changed(filename: str) -> None:
102 if current_worker["instance"]:
103 current_worker["instance"].shutdown()
104
105 def signal_shutdown(signalnum: int, _: Any) -> None:
106 should_restart["value"] = False
107 if current_worker["instance"]:
108 current_worker["instance"].shutdown()
109
110 signal.signal(signal.SIGTERM, signal_shutdown)
111 signal.signal(signal.SIGINT, signal_shutdown)
112
113 reloader = Reloader(callback=file_changed, watch_html=False)
114 reloader.start()
115
116 while should_restart["value"]:
117 w = Worker(**worker_kwargs)
118 current_worker["instance"] = w
119 w.run()
120
121
122def _run_once(worker_kwargs: dict[str, Any]) -> None:
123 w = Worker(**worker_kwargs)
124
125 def _shutdown(signalnum: int, _: Any) -> None:
126 logger.info(
127 "Job worker shutdown signal received",
128 extra={"signalnum": signalnum},
129 )
130 w.shutdown()
131
132 signal.signal(signal.SIGTERM, _shutdown)
133 signal.signal(signal.SIGINT, _shutdown)
134
135 w.run()
136
137
138@cli.command()
139def clear() -> None:
140 """Clear completed job results"""
141 cutoff = timezone.now() - datetime.timedelta(
142 seconds=settings.JOBS_RESULTS_RETENTION
143 )
144 click.echo(f"Clearing job results created before {cutoff}")
145 count = JobResult.query.filter(created_at__lt=cutoff).delete()
146 click.echo(f"Deleted {count} jobs")
147
148
149@cli.command()
150def stats() -> None:
151 """Show job queue statistics"""
152 pending = JobRequest.query.count()
153 processing = JobProcess.query.count()
154
155 successful = JobResult.query.successful().count()
156 errored = JobResult.query.errored().count()
157 lost = JobResult.query.lost().count()
158
159 click.secho(f"Pending: {pending}", bold=True)
160 click.secho(f"Processing: {processing}", bold=True)
161 click.secho(f"Successful: {successful}", bold=True, fg="green")
162 click.secho(f"Errored: {errored}", bold=True, fg="red")
163 click.secho(f"Lost: {lost}", bold=True, fg="yellow")
164
165
166@cli.command()
167@click.option("--yes", "-y", is_flag=True, help="Skip confirmation prompt.")
168def purge(yes: bool) -> None:
169 """Delete all pending and running jobs"""
170 if not yes and not click.confirm(
171 "Are you sure you want to clear all running and pending jobs? This will delete all current Jobs and JobRequests"
172 ):
173 return
174
175 deleted = JobRequest.query.all().delete()
176 click.echo(f"Deleted {deleted} job requests")
177
178 deleted = JobProcess.query.all().delete()
179 click.echo(f"Deleted {deleted} jobs")
180
181
182@cli.command()
183@click.argument("job_class_name", type=str)
184def run(job_class_name: str) -> None:
185 """Run a job directly without a worker"""
186 job = jobs_registry.load_job(job_class_name, {"args": [], "kwargs": {}})
187 click.secho("Loaded job: ", bold=True, nl=False)
188 print(job)
189 job.run()
190
191
192@cli.command("list")
193def list_jobs() -> None:
194 """List all registered jobs"""
195 for name, job_class in jobs_registry.jobs.items():
196 click.secho(name, bold=True, nl=False)
197 description = job_class.__doc__.strip() if job_class.__doc__ else ""
198 if description:
199 click.secho(f": {description}", dim=True)
200 else:
201 click.echo("")