v0.163.0
  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("")