Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
30 commits
Select commit Hold shift + click to select a range
e6c6e9f
feat(jobs): describe creatable job types and validate job settings on…
mihow Sep 30, 2026
b40050f
feat(ui): add job type discovery and schema-to-form mapping helpers
mihow Sep 30, 2026
06d9bb8
docs: explain how a job type appears in the generated Create Job dial…
mihow Sep 30, 2026
b2ae0fc
feat(ui): add Create job strings and server error mapping
mihow Sep 30, 2026
992e30a
feat(ui): replace the New job dialog with a schema-driven Create job …
mihow Sep 30, 2026
06f9739
fix(ml): allow filtering the algorithms list by task type
mihow Sep 30, 2026
c801b55
fix(ui): pad the Create job form and name jobs without the capture count
mihow Sep 30, 2026
7905681
fix(jobs): store no params for job types that take none
mihow Sep 30, 2026
8dbbb2d
Merge main into feat/jobs-panel
mihow Sep 30, 2026
337b082
fix(jobs): keep accepting ML jobs created without a pipeline
mihow Sep 30, 2026
225ca12
feat(jobs): mark staff-only settings and tidy generated labels and he…
mihow Sep 30, 2026
20c8672
feat(ui): tuck staff settings away, start jobs by default, and show u…
mihow Sep 30, 2026
ecac37c
fix(jobs): say "Only a superuser can change this setting." like the t…
mihow Sep 30, 2026
47a4c0d
refactor(jobs): gate post-processing methods by project feature flag …
mihow Oct 2, 2026
905ed72
refactor(ui): render job settings straight from the pydantic schema
mihow Oct 2, 2026
b316881
docs: update the jobs panel reference for feature flags and plain pyd…
mihow Oct 2, 2026
58ee1f9
refactor(jobs): serve a settings model's schema with no wrapper
mihow Oct 2, 2026
2b5766f
refactor(jobs): describe every job type's inputs with one pydantic mo…
mihow Oct 2, 2026
15aeca9
refactor(ui): send every job input in params and render one schema pe…
mihow Oct 2, 2026
1f93472
fix(post-processing): set the capture set column on jobs started from…
mihow Oct 2, 2026
d840a17
fix(ui): show the settings divider only for a method's settings
mihow Oct 2, 2026
163d5c1
feat(jobs): require a job's pipeline to be enabled for its project, a…
mihow Oct 2, 2026
dbb3c83
fix(ui): name a new job after its capture set when one is picked
mihow Oct 2, 2026
10a133f
fix(jobs): accept only pipelines the project has switched on
mihow Oct 2, 2026
510ec38
fix(ui): type picker rows, translate number validation, and recover a…
mihow Oct 2, 2026
abfa352
fix(jobs): refuse changing a job's project or type, check algorithms …
mihow Oct 2, 2026
bc0c81b
fix(ui): give the dialog's pickers, selects and checkboxes accessible…
mihow Oct 2, 2026
dd6f8e4
feat(jobs): name job choices as actions and group them by what the us…
mihow Oct 6, 2026
95f896f
feat(ui): pick a job from one grouped list and show friendly job names
mihow Oct 6, 2026
9fdffb2
test(jobs): build the job types fixtures once per class
mihow Oct 6, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
273 changes: 272 additions & 1 deletion ami/jobs/models.py
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import datetime
import inspect
import logging
import random
import time
Expand All @@ -11,15 +12,27 @@
from django.conf import settings
from django.db import models, transaction
from django.utils.text import slugify
from django.utils.translation import gettext_lazy as _
from django_pydantic_field import SchemaField
from guardian.shortcuts import get_perms
from rest_framework import serializers

from ami.base.models import BaseModel
from ami.base.schemas import ConfigurableStage, ConfigurableStageParam
from ami.jobs.schemas import (
JOB_GROUP_LABELS,
CaptureSetJobConfig,
JobGroup,
JobGroupDescription,
JobTypeDescription,
JobTypeVariantDescription,
MLJobConfig,
StationJobConfig,
)
from ami.jobs.tasks import cleanup_async_job_if_needed, run_job
from ami.main.models import Deployment, Project, SourceImage, SourceImageCollection
from ami.ml.models import Pipeline
from ami.ml.post_processing.registry import get_postprocessing_task
from ami.ml.post_processing.registry import POSTPROCESSING_TASKS, get_postprocessing_task
from ami.utils.schemas import OrderedEnum

logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -433,6 +446,88 @@ def emit(self, record: logging.LogRecord):
logger.error(f"Failed to save log for job #{self.job.pk}: {e}")


# Config fields that are also stored on a Job column, so the jobs list can filter and join on them.
JOB_COLUMNS = ("pipeline_id", "source_image_collection_id", "source_image_single_id", "deployment_id")


def entity_fields(model: type[pydantic.BaseModel]) -> dict[str, str]:
"""Map each config field carrying an ``ami_entity`` hint to that entity (an API route)."""
return {
name: prop["ami_entity"]
for name, prop in model.schema().get("properties", {}).items()
if isinstance(prop, dict) and "ami_entity" in prop
}


def pydantic_messages(exc: pydantic.ValidationError) -> list[str]:
"""Flatten a pydantic error into ``"field: message"`` lines for a 400 response."""
messages = []
for err in exc.errors():
field = ".".join(str(part) for part in err.get("loc", ()) if part != "__root__")
messages.append(f"{field}: {err['msg']}" if field else err["msg"])
return messages


def _entity_queryset(entity: str, project: Project | None):
"""The rows of ``entity`` a job in ``project`` may refer to, or None when not project-scoped."""
from django.db.models import Q

from ami.main.models import Event, Occurrence, TaxaList
from ami.ml.models import Algorithm

scoped = {
"captures/collections": lambda: SourceImageCollection.objects.filter(project=project),
"deployments": lambda: Deployment.objects.filter(project=project),
"captures": lambda: SourceImage.objects.filter(project=project),
# Algorithms are shared catalogue rows: any existing one may be named.
"ml/algorithms": lambda: Algorithm.objects.all(),
# A pipeline is shared; a job may use the ones its project has enabled.
"ml/pipelines": lambda: Pipeline.objects.filter(
project_pipeline_configs__project=project, project_pipeline_configs__enabled=True
),
"events": lambda: Event.objects.filter(project=project),
"occurrences": lambda: Occurrence.objects.filter(project=project),
# Public lists belong to no project and may be used by any.
"taxa/lists": lambda: TaxaList.objects.filter(Q(projects=project) | Q(projects__isnull=True)),
}
Comment thread
mihow marked this conversation as resolved.
factory = scoped.get(entity)
return factory() if factory else None


def check_entities_in_project(values: dict, entities: dict[str, str], project: Project | None) -> None:
"""Refuse ids in ``values`` that point outside ``project``.

``entities`` maps a field name to its API entity. The schema can only say an id is an
integer; this is the check that it names a row the job's project owns.
"""
errors = []
for field, entity in entities.items():
value = values.get(field)
if value in (None, [], ""):
continue
ids = set(value) if isinstance(value, (list, tuple)) else {value}
queryset = _entity_queryset(entity, project)
if queryset is None:
continue
found = set(queryset.filter(pk__in=ids).values_list("pk", flat=True).distinct())
missing = sorted(ids - found)
if missing:
errors.append(f"{field}: {missing} not found in this project.")
if errors:
raise serializers.ValidationError({"params": {"config": errors}})
Comment thread
coderabbitai[bot] marked this conversation as resolved.


def _validate_config(model_cls: type[pydantic.BaseModel], config, project: Project | None) -> pydantic.BaseModel:
if not isinstance(config, dict):
raise serializers.ValidationError({"params": {"config": "Must be an object."}})
try:
model = model_cls(**config)
except pydantic.ValidationError as exc:
raise serializers.ValidationError({"params": {"config": pydantic_messages(exc)}})
check_entities_in_project(model.dict(), entity_fields(model_cls), project)
return model


@dataclass
class JobType:
"""
Expand All @@ -441,8 +536,83 @@ class JobType:
Job types must be defined as classes because they define code, not just configuration.
"""

# A fixed internal name, used in logs and stage names. Users see ``label``.
name: str
key: str
# What users see for this job type in the Create Job picker and the jobs list, as an action
# ("Process captures"). Wrap it in gettext_lazy to translate it; left empty, ``name`` is used.
label: str = ""
# The Create Job picker heading this type is listed under. A type with variants leaves it
# empty and each variant declares its own.
group: JobGroup | None = None
# Help text under the job type select. Wrap it in gettext_lazy to translate it; left empty,
# the first paragraph of the class docstring is used.
description: str = ""

# Whether a person can start one from the Create Job dialog. The rest are created by
# the platform for the user: an export from the exports page, for example.
user_creatable: bool = False

# Everything a new job of this type takes, as a pydantic model: the Create Job dialog renders
# it and the API validates against it. See ami/jobs/schemas.py.
config_schema: type[pydantic.BaseModel] | None = None

# A job type whose work is chosen from a registry (post-processing tasks) names the
# ``params`` key that holds the choice, and lists the choices as variants.
variant_key: str | None = None

@classmethod
def help_text(cls) -> str:
"""``description``, else the first paragraph of the class docstring."""
if cls.description:
return str(cls.description)
return (inspect.getdoc(cls) or "").split("\n\n")[0].replace("\n", " ").strip()

@classmethod
def label_for(cls, params: dict | None) -> str:
"""The label users see for one job of this type."""
return str(cls.label or cls.name)

@classmethod
def variants(cls, project: Project) -> list[JobTypeVariantDescription]:
return []

@classmethod
def describe(cls, project: Project, allowed: bool) -> JobTypeDescription | None:
"""What the Create Job dialog shows for this type, or None when it has nothing to offer."""
variants = cls.variants(project)
if cls.variant_key and not variants:
return None # e.g. post-processing with no method turned on for this project
return JobTypeDescription(
key=cls.key,
name=str(cls.label or cls.name),
description=cls.help_text(),
group=cls.group,
allowed=allowed,
config_schema=cls.config_schema.schema() if cls.config_schema else None,
variant_key=cls.variant_key,
variants=variants,
)

@classmethod
def column_ids(cls, params: dict) -> dict[str, int]:
"""The Job column ids (``pipeline_id``, ...) carried in validated params."""
config = params.get("config") or {}
return {field: config[field] for field in JOB_COLUMNS if config.get(field) is not None}

@classmethod
def validate_params(cls, project: Project | None, user, params) -> dict:
"""Check a new job's ``params`` before it is saved and return what should be stored.

Raises ``serializers.ValidationError`` (a 400) for a bad value. The config is validated
against ``config_schema`` and every id in it must belong to the project.
"""
if not isinstance(params, dict):
raise serializers.ValidationError({"params": "Must be an object."})
if cls.config_schema is None:
return {}
model = _validate_config(cls.config_schema, params.get("config") or {}, project)
return {"config": model.dict()}

# @TODO Consider adding custom vocabulary for job types to be used in the UI
# verb: str = "Sync"
Expand All @@ -460,6 +630,11 @@ def run(cls, job: "Job"):
class MLJob(JobType):
name = "ML pipeline"
key = "ml"
label = _("Process captures")
group = JobGroup.PROCESS_IMAGES
description = _("Detects insects and predicts their species with a pipeline.")
user_creatable = True
config_schema = MLJobConfig

@classmethod
def run(cls, job: "Job"):
Expand Down Expand Up @@ -710,6 +885,11 @@ class DataStorageSyncJob(JobType):

name = "Data storage sync"
key = "data_storage_sync"
label = _("Sync captures from storage")
group = JobGroup.ORGANIZE_CAPTURES
description = _("Add new captures from a station's data storage, then regroup them into sessions.")
user_creatable = True
config_schema = StationJobConfig
regroup_stage_key = "regroup_sessions"
regroup_stage_name = "Regroup sessions"

Expand Down Expand Up @@ -804,6 +984,11 @@ def run(cls, job: "Job"):
class SourceImageCollectionPopulateJob(JobType):
name = "Populate capture set"
key = "populate_captures_collection"
label = _("Fill a capture set")
group = JobGroup.ORGANIZE_CAPTURES
description = _("Fill a capture set with the captures its sampling method selects.")
user_creatable = True
config_schema = CaptureSetJobConfig

@classmethod
def run(cls, job: "Job"):
Expand Down Expand Up @@ -891,6 +1076,61 @@ def run(cls, job: "Job"):
class PostProcessingJob(JobType):
name = "Post Processing"
key = "post_processing"
description = _(
"Revise existing results with a post-processing method, such as masking classes or "
"filtering out detections too small to identify."
)
user_creatable = True
variant_key = "task"

@classmethod
def enabled_tasks(cls, project: Project | None) -> dict:
"""The registered tasks whose feature flag is on for ``project``; the others are hidden."""
flags = project.feature_flags if project else None
return {key: task for key, task in POSTPROCESSING_TASKS.items() if flags and getattr(flags, task.feature_flag)}

@classmethod
def label_for(cls, params: dict | None) -> str:
"""A post-processing job is named by its method, e.g. "Mark detections too small to identify"."""
task_key = (params or {}).get("task")
task_cls = get_postprocessing_task(task_key) if isinstance(task_key, str) else None
return str(task_cls.label) if task_cls else super().label_for(params)

@classmethod
def variants(cls, project: Project) -> list[JobTypeVariantDescription]:
return [
JobTypeVariantDescription(
key=key,
name=str(task_cls.label),
description=str(task_cls.description),
group=task_cls.group,
config_schema=task_cls.config_schema.schema(),
)
for key, task_cls in cls.enabled_tasks(project).items()
]

@classmethod
def validate_params(cls, project: Project | None, user, params) -> dict:
"""Check a post-processing job's ``{"task": ..., "config": {...}}`` before it is saved.

The task must be turned on for the project (its feature flag), the config must pass the
task's schema, and every id in it must belong to the project. Returns the params with the
config normalized by the schema, so the stored job carries every default the worker uses.
"""
if not isinstance(params, dict) or set(params) - {"task", "config"}:
raise serializers.ValidationError(
{"params": 'Post-processing jobs take params of the form {"task": <key>, "config": {...}}.'}
)
task_key = params.get("task")
task_cls = get_postprocessing_task(task_key) if isinstance(task_key, str) else None
if task_cls is None:
raise serializers.ValidationError({"params": {"task": f"Unknown post-processing task {task_key!r}."}})
if task_key not in cls.enabled_tasks(project):
raise serializers.ValidationError(
{"params": {"task": f"{task_cls.name} is not turned on for this project."}}
)
model = _validate_config(task_cls.config_schema, params.get("config") or {}, project)
return {"task": task_key, "config": model.dict()}

@classmethod
def run(cls, job: "Job"):
Expand Down Expand Up @@ -940,6 +1180,11 @@ class RegroupEventsJob(JobType):

name = "Regroup sessions"
key = "regroup_events"
label = _("Regroup captures into sessions")
group = JobGroup.ORGANIZE_CAPTURES
description = _("Regroup a station's captures into sessions using the project's session time gap.")
user_creatable = True
config_schema = StationJobConfig

@classmethod
def run(cls, job: "Job"):
Expand Down Expand Up @@ -988,6 +1233,27 @@ def run(cls, job: "Job"):
]


def describe_job_groups(job_types: list[JobTypeDescription]) -> list[JobGroupDescription]:
"""The picker headings ``job_types`` use, in ``JobGroup`` order."""
used = {jt.group for jt in job_types} | {v.group for jt in job_types for v in jt.variants}
return [JobGroupDescription(key=group, label=str(JOB_GROUP_LABELS[group])) for group in JobGroup if group in used]


def describe_job_types(project: Project, user) -> list[JobTypeDescription]:
"""The job types ``user`` may pick in the Create Job dialog for ``project``.

A type the user may not run is still listed with ``allowed=False``, so the dialog can show it
disabled. Permissions are read once for the whole list.
"""
perms = set(get_perms(user, project))
described = (
job_type.describe(project, allowed=user.is_superuser or f"run_{job_type.key}_job" in perms)
for job_type in VALID_JOB_TYPES
if job_type.user_creatable
)
return [description for description in described if description]


def get_job_type_by_key(key: str) -> type[JobType] | None:
for job_type in VALID_JOB_TYPES:
if job_type.key == key:
Expand Down Expand Up @@ -1395,6 +1661,11 @@ def check_custom_permission(self, user, action: str) -> bool:
permission_codename = f"{action}_{job_type}_job"

project = self.get_project() if hasattr(self, "get_project") else None
if job_type == PostProcessingJob.key and action in ("run", "retry") and not user.is_superuser:
# Turning a method's feature flag off also stops its existing jobs being re-run.
task_key = (self.params or {}).get("task")
if task_key not in PostProcessingJob.enabled_tasks(project):
return False
return user.has_perm(permission_codename, project)

def get_custom_user_permissions(self, user) -> list[str]:
Expand Down
Loading