From 610bd6582196b186d8e6842e4f8b3a7b51134922 Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Mon, 5 Oct 2026 14:56:38 -0700 Subject: [PATCH 01/12] feat: record which job created each detection and classification Detections and classifications gain a nullable `job` foreign key. Pipeline result saving, class masking and the size filter set it on the rows they create; rows that already exist keep the job they had. The detection and classification API responses expose it as a read-only id. Co-Authored-By: Claude Opus 5.5 (1M context) --- ami/main/api/serializers.py | 4 ++ .../0096_detection_and_classification_job.py | 36 +++++++++++++++++ ami/main/models.py | 11 +++-- ami/ml/models/pipeline.py | 12 ++++++ ami/ml/post_processing/class_masking.py | 9 +++++ ami/ml/post_processing/small_size_filter.py | 1 + .../tests/test_class_masking.py | 24 +++++++++++ ami/ml/tests.py | 40 +++++++++++++++++++ 8 files changed, 131 insertions(+), 6 deletions(-) create mode 100644 ami/main/migrations/0096_detection_and_classification_job.py diff --git a/ami/main/api/serializers.py b/ami/main/api/serializers.py index 962c802af..1ca588980 100644 --- a/ami/main/api/serializers.py +++ b/ami/main/api/serializers.py @@ -1105,6 +1105,7 @@ class ClassificationSerializer(DefaultSerializer): algorithm = AlgorithmSerializer(read_only=True) top_n = ClassificationPredictionItemSerializer(many=True, read_only=True) applied_to = ClassificationAppliedToSerializer(read_only=True) + job = serializers.PrimaryKeyRelatedField(read_only=True, help_text="The job that wrote this classification.") class Meta: model = Classification @@ -1118,6 +1119,7 @@ class Meta: "logits", "top_n", "applied_to", + "job", "created_at", "updated_at", ] @@ -1263,6 +1265,7 @@ class DetectionSerializer(DefaultSerializer): queryset=Algorithm.objects.all(), source="detection_algorithm", write_only=True ) classifications = ClassificationNestedSerializer(many=True, read_only=True) + job = serializers.PrimaryKeyRelatedField(read_only=True, help_text="The job that wrote this detection.") class Meta: model = Detection @@ -1271,6 +1274,7 @@ class Meta: "detection_algorithm", "detection_algorithm_id", "classifications", + "job", ] diff --git a/ami/main/migrations/0096_detection_and_classification_job.py b/ami/main/migrations/0096_detection_and_classification_job.py new file mode 100644 index 000000000..21c25e720 --- /dev/null +++ b/ami/main/migrations/0096_detection_and_classification_job.py @@ -0,0 +1,36 @@ +# Generated by Django 4.2.10 on 2026-10-05 17:41 + +from django.db import migrations, models +import django.db.models.deletion + + +class Migration(migrations.Migration): + dependencies = [ + ("jobs", "0023_alter_job_job_type_key"), + ("main", "0095_grant_sync_deployment_to_mldatamanager"), + ] + + operations = [ + migrations.AddField( + model_name="classification", + name="job", + field=models.ForeignKey( + blank=True, + null=True, + on_delete=django.db.models.deletion.SET_NULL, + related_name="classifications", + to="jobs.job", + ), + ), + migrations.AddField( + model_name="detection", + name="job", + field=models.ForeignKey( + blank=True, + null=True, + on_delete=django.db.models.deletion.SET_NULL, + related_name="detections", + to="jobs.job", + ), + ), + ] diff --git a/ami/main/models.py b/ami/main/models.py index 9943c111a..24207ed63 100644 --- a/ami/main/models.py +++ b/ami/main/models.py @@ -2995,7 +2995,9 @@ class Classification(BaseModel): null=True, related_name="classifications", ) - # job = models.CharField(max_length=255, null=True) + job = models.ForeignKey( + "jobs.Job", on_delete=models.SET_NULL, null=True, blank=True, related_name="classifications" + ) applied_to = models.ForeignKey( "self", on_delete=models.SET_NULL, @@ -3189,11 +3191,8 @@ class Detection(BaseModel): # @TODO not sure if this detection score is ever used # I think it was intended to be the score of the detection algorithm (bbox score) detection_score = models.FloatField(null=True, blank=True) - # detection_job = models.ForeignKey( - # "Job", - # on_delete=models.SET_NULL, - # null=True, - # ) + # The run that wrote this detection, when there was one. + job = models.ForeignKey("jobs.Job", on_delete=models.SET_NULL, null=True, blank=True, related_name="detections") similarity_vector = models.JSONField(null=True, blank=True) diff --git a/ami/ml/models/pipeline.py b/ami/ml/models/pipeline.py index 8ebaf6d6d..3d5b77d58 100644 --- a/ami/ml/models/pipeline.py +++ b/ami/ml/models/pipeline.py @@ -523,6 +523,7 @@ def get_or_create_detection( algorithms_known: dict[str, Algorithm], save: bool = True, logger: logging.Logger = logger, + job_id: int | None = None, ) -> tuple[Detection, bool]: """ Create a Detection object from a DetectionResponse, or update an existing one. @@ -616,6 +617,7 @@ def get_or_create_detection( path=crop_url, detection_time=detection_resp.timestamp, detection_algorithm=detection_algo, + job_id=job_id, ) if save: new_detection.save() @@ -633,6 +635,7 @@ def create_detections( detections: list[DetectionResponse], algorithms_known: dict[str, Algorithm], logger: logging.Logger = logger, + job_id: int | None = None, ) -> list[Detection]: """ Efficiently create multiple Detection objects from a list of DetectionResponse objects, grouped by source image. @@ -663,6 +666,7 @@ def create_detections( algorithms_known=algorithms_known, save=False, logger=logger, + job_id=job_id, ) if created: new_detections.append(detection) @@ -756,6 +760,7 @@ def create_classification( algorithms_known: dict[str, Algorithm], save: bool = True, logger: logging.Logger = logger, + job_id: int | None = None, ) -> tuple[Classification, bool]: """ Create a Classification object from a ClassificationResponse, or update an existing one. @@ -843,6 +848,7 @@ def create_classification( scores=classification_resp.scores, terminal=classification_resp.terminal, category_map=classification_algo.category_map, + job_id=job_id, ) classification = new_classification @@ -861,6 +867,7 @@ def create_classifications( algorithms_known: dict[str, Algorithm], logger: logging.Logger = logger, save: bool = True, + job_id: int | None = None, ) -> list[Classification]: """ Efficiently create multiple Classification objects from a list of ClassificationResponse objects, @@ -885,6 +892,7 @@ def create_classifications( algorithms_known=algorithms_known, save=False, logger=logger, + job_id=job_id, ) if created: new_classifications.append(classification) @@ -1071,10 +1079,12 @@ def save_results( "Algorithms and category maps must be registered before processing, using /info endpoint." ) + # New rows record the job that wrote them; rows that already existed keep theirs. detections = create_detections( detections=results.detections, algorithms_known=algorithms_known, logger=job_logger, + job_id=job.pk if job else None, ) classifications = create_classifications( @@ -1082,6 +1092,7 @@ def save_results( detection_responses=results.detections, algorithms_known=algorithms_known, logger=job_logger, + job_id=job.pk if job else None, ) # Create a new occurrence for each detection (no tracking yet) @@ -1123,6 +1134,7 @@ def save_results( detections=null_detection_responses, algorithms_known=algorithms_known, logger=job_logger, + job_id=job.pk if job else None, ) total_time = time.time() - start_time diff --git a/ami/ml/post_processing/class_masking.py b/ami/ml/post_processing/class_masking.py index 2da2001b7..18b22c26d 100644 --- a/ami/ml/post_processing/class_masking.py +++ b/ami/ml/post_processing/class_masking.py @@ -1,4 +1,5 @@ import logging +import typing from collections.abc import Callable import numpy as np @@ -11,6 +12,9 @@ from ami.ml.models.algorithm import Algorithm, AlgorithmTaskType from ami.ml.post_processing.base import BasePostProcessingTask +if typing.TYPE_CHECKING: + from ami.jobs.models import Job + logger = logging.getLogger(__name__) @@ -52,6 +56,7 @@ def make_classifications_filtered_by_taxa_list( task_logger: logging.Logger = logger, on_setup: Callable[[int], None] | None = None, on_batch: Callable[[dict], None] | None = None, + job: "Job | None" = None, ) -> dict[str, int]: """Re-score ``classifications`` by masking out classes absent from ``taxa_list``. @@ -74,6 +79,8 @@ def make_classifications_filtered_by_taxa_list( ``occurrences_updated`` counts only occurrences whose determination actually changed (not just any occurrence touched), matching the size-filter convention. + New classifications record ``job`` as the run that wrote them. + Returns final counters (checked / masked / occurrences updated) for stage metrics. """ taxa_in_list = set(taxa_list.taxa.all()) @@ -180,6 +187,7 @@ def make_classifications_filtered_by_taxa_list( terminal=True, timestamp=classification.timestamp, applied_to=classification, + job=job, created_at=timestamp, updated_at=timestamp, ) @@ -368,6 +376,7 @@ def _on_batch(m: dict) -> None: task_logger=self.logger, on_setup=_on_setup, on_batch=_on_batch, + job=self.job, ) self.report_stage_metrics(metrics) self.logger.info(f"=== Completed {self.name} ===") diff --git a/ami/ml/post_processing/small_size_filter.py b/ami/ml/post_processing/small_size_filter.py index 6a39af780..f0b8289eb 100644 --- a/ami/ml/post_processing/small_size_filter.py +++ b/ami/ml/post_processing/small_size_filter.py @@ -126,6 +126,7 @@ def run(self) -> None: timestamp=timezone.now(), # How is this different from created_at? algorithm=self.algorithm, applied_to=None, # Size filter is applied to original detection, not a previous classification + job=self.job, ) ) detections_to_update.add(det) diff --git a/ami/ml/post_processing/tests/test_class_masking.py b/ami/ml/post_processing/tests/test_class_masking.py index bf38decc9..541a70195 100644 --- a/ami/ml/post_processing/tests/test_class_masking.py +++ b/ami/ml/post_processing/tests/test_class_masking.py @@ -253,6 +253,30 @@ def test_task_run_collection_scope_persists_masking_algorithm(self): occ.refresh_from_db() self.assertEqual(occ.determination, self.species_taxa[1], "Occurrence determination follows the masked result") + def test_task_run_records_the_job_on_new_classifications_only(self): + from ami.jobs.models import Job + + logits = [0.5, 3.0, 3.5] + taxa_list = TaxaList.objects.create(name="Regional list") + taxa_list.taxa.set(self.species_taxa[:2]) + det, _ = self._detection_with_occurrence() + original = self._create_classification_with_logits(det, self.species_taxa[2], _softmax(logits), logits) + job = Job.objects.create(project=self.project, name="Class masking job test", job_type_key="post_processing") + job.progress.add_stage("Post Processing", key="post_processing") + job.save() + + ClassMaskingTask( + job=job, + source_image_collection_id=self.collection.pk, + taxa_list_id=taxa_list.pk, + algorithm_id=self.algorithm.pk, + ).run() + + new_clf = Classification.objects.get(detection=det, terminal=True, applied_to=original) + self.assertEqual(new_clf.job_id, job.pk) + original.refresh_from_db() + self.assertIsNone(original.job_id, "The source classification keeps the job it had") + def test_rerun_does_not_duplicate_masked_classifications(self): """Re-running the same mask must not create a second masked classification for a source already re-scored, even if that source became terminal again in between. diff --git a/ami/ml/tests.py b/ami/ml/tests.py index bd92bb02f..b5f99a66e 100644 --- a/ami/ml/tests.py +++ b/ami/ml/tests.py @@ -672,6 +672,24 @@ def test_save_results(self): # @TODO test the cached counts for detections, etc are updated on Events, Deployments, etc. + def test_save_results_records_the_job_on_the_rows_it_creates(self): + from ami.jobs.models import Job + + first_job = Job.objects.create(project=self.project, name="First run", pipeline=self.pipeline) + second_job = Job.objects.create(project=self.project, name="Second run", pipeline=self.pipeline) + results = self.fake_pipeline_results(self.test_images, self.pipeline) + + save_results(results, job_id=first_job.pk) + # A second run returning the same boxes and labels reuses the rows, which keep the first job. + save_results(results, job_id=second_job.pk) + + detections = Detection.objects.filter(source_image__in=self.test_images) + classifications = Classification.objects.filter(detection__in=detections) + self.assertTrue(detections.exists()) + self.assertTrue(classifications.exists()) + self.assertEqual(set(detections.values_list("job_id", flat=True)), {first_job.pk}) + self.assertEqual(set(classifications.values_list("job_id", flat=True)), {first_job.pk}) + def test_skip_existing_when_all_matching(self): """ When processing images, skip images that have already been processed by the same set of algorithms. @@ -1603,6 +1621,28 @@ def test_run_reports_stage_metrics_on_job(self): # count equals the detection count. self.assertEqual(params.get("occurrences_updated"), total) + def test_new_classifications_record_the_job(self): + """The size filter's "Not identifiable" classifications point at the job that ran it.""" + from ami.jobs.models import Job + + for image in self.collection.images.all(): + Detection.objects.create( + source_image=image, + bbox=[0, 0, 10, 10], # small → flagged + created_at=datetime.datetime.now(datetime.timezone.utc), + ).associate_new_occurrence() + job = Job.objects.create(project=self.project, name="size filter job test", job_type_key="post_processing") + job.progress.add_stage("Post Processing", key="post_processing") + job.save() + + SmallSizeFilterTask(job=job, source_image_collection_id=self.collection.pk, size_threshold=0.01).run() + + new_classifications = Classification.objects.filter( + detection__source_image__in=self.collection.images.all(), taxon__name="Not identifiable" + ) + self.assertTrue(new_classifications.exists()) + self.assertEqual(set(new_classifications.values_list("job_id", flat=True)), {job.pk}) + def test_progress_save_bumps_updated_at_for_reaper(self): """A progress heartbeat bumps ``Job.updated_at`` so the stale-job reaper leaves an actively-running post-processing job alone. From 4e0a34716803fa45c5ffc377329057cbe44fcfa5 Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Mon, 5 Oct 2026 14:56:41 -0700 Subject: [PATCH 02/12] feat: filter occurrences by the job that wrote their results `?job=` on the occurrence list matches occurrences with a detection or a classification written by that job, using EXISTS subqueries so each occurrence appears once. A non-integer id returns 400. Co-Authored-By: Claude Opus 5.5 (1M context) --- ami/main/api/views.py | 25 ++++++++++ ami/main/models.py | 11 +++++ ami/main/tests.py | 106 ++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 142 insertions(+) diff --git a/ami/main/api/views.py b/ami/main/api/views.py index fdc72c963..84747bb59 100644 --- a/ami/main/api/views.py +++ b/ami/main/api/views.py @@ -1448,10 +1448,29 @@ def filter_queryset(self, request, queryset, view): return queryset +class OccurrenceJobFilter(filters.BaseFilterBackend): + """ + Filter occurrences by the job that wrote one of their detections or classifications. + """ + + query_param = "job" + + def filter_queryset(self, request, queryset, view): + job_id = SingleParamSerializer[int].clean( + param_name=self.query_param, + field=serializers.IntegerField(required=False, min_value=1), + data=request.query_params, + ) + if job_id is None: + return queryset + return queryset.written_by_job(job_id) + + OCCURRENCE_FILTER_BACKENDS = ( CustomOccurrenceDeterminationFilter, OccurrenceCollectionFilter, OccurrenceAlgorithmFilter, + OccurrenceJobFilter, OccurrenceDateFilter, OccurrenceVerified, OccurrenceVerifiedByMeFilter, @@ -1558,6 +1577,12 @@ def get_queryset(self) -> QuerySet["Occurrence"]: required=False, type=OpenApiTypes.INT, ), + OpenApiParameter( + name="job", + description="Filter occurrences by the job that wrote one of their detections or classifications.", + required=False, + type=OpenApiTypes.INT, + ), ] ) def list(self, request, *args, **kwargs): diff --git a/ami/main/models.py b/ami/main/models.py index 24207ed63..0f6c6686e 100644 --- a/ami/main/models.py +++ b/ami/main/models.py @@ -3387,6 +3387,17 @@ def not_processed_by_algorithm(self, algorithm_ids) -> "OccurrenceQuerySet": """Occurrences with no result from any of the given algorithms.""" return self.exclude(self._processed_by_algorithm_q(algorithm_ids)) + def written_by_job(self, job_id: int) -> "OccurrenceQuerySet": + """Occurrences with a detection or a classification written by the given job. + + Identifications are not matched: people make them, not jobs. Two EXISTS subqueries + return each occurrence once, where a join would return one row per matching result. + """ + return self.filter( + Exists(Detection.objects.filter(occurrence=OuterRef("pk"), job_id=job_id)) + | Exists(Classification.objects.filter(detection__occurrence=OuterRef("pk"), job_id=job_id)) + ) + def with_timestamps(self): """ These are timestamps used for filtering and ordering in the UI. diff --git a/ami/main/tests.py b/ami/main/tests.py index aab6d943d..de82178c6 100644 --- a/ami/main/tests.py +++ b/ami/main/tests.py @@ -7242,6 +7242,112 @@ def test_exclude_is_the_complement_of_include(self): self.assertEqual(included & excluded, set()) +class TestOccurrenceJobFilter(APITestCase): + """ + Covers the ``?job=`` occurrence filter: occurrences with a detection or a + classification written by the given job. Each occurrence must appear once, however + many of its rows the job wrote, and a malformed id must be a 400, not a 500. + """ + + def setUp(self): + from ami.main.models import Taxon + from ami.ml.models.algorithm import Algorithm + + self.project = Project.objects.create(name="Occurrence Job Filter Project") + self.deployment = Deployment.objects.create(project=self.project, name="dep") + self.event = Event.objects.create( + project=self.project, + deployment=self.deployment, + group_by="2024-01-01", + start=datetime.datetime(2024, 1, 1, 0, 0), + ) + self.source_image = SourceImage.objects.create( + deployment=self.deployment, + project=self.project, + event=self.event, + path="occ-job-filter.jpg", + ) + self.taxon = Taxon.objects.create(name="Occurrence Job Filter Taxon") + self.algorithm = Algorithm.objects.create(name="Job Filter Algorithm", version=1, task_type="classification") + self.job = Job.objects.create(project=self.project, name="Job under test") + self.other_job = Job.objects.create(project=self.project, name="Other job") + + # The job wrote the detection only. + self.occ_detected = self._make_occurrence([(self.job, [None])]) + # The job wrote a classification on another job's detection. + self.occ_classified = self._make_occurrence([(self.other_job, [self.job])]) + # The job wrote three detections and their classifications: the duplicate-row trap. + self.occ_multi = self._make_occurrence([(self.job, [self.job, self.job])] * 3) + # Only the other job, and no job at all, must not match. + self.occ_other = self._make_occurrence([(self.other_job, [self.other_job])]) + self.occ_none = self._make_occurrence([(None, [None])]) + + def _make_occurrence(self, detections) -> Occurrence: + """``detections`` is a list of (detection job, [classification jobs]).""" + occ = Occurrence.objects.create( + project=self.project, + event=self.event, + deployment=self.deployment, + determination=self.taxon, + determination_score=0.9, + ) + for detection_job, classification_jobs in detections: + detection = Detection.objects.create( + source_image=self.source_image, + bbox=[0.0, 0.0, 1.0, 1.0], + occurrence=occ, + job=detection_job, + ) + for classification_job in classification_jobs: + detection.classifications.create( + taxon=self.taxon, + algorithm=self.algorithm, + score=0.9, + timestamp=datetime.datetime.now(), + job=classification_job, + ) + return occ + + def _list(self, job_param): + return self.client.get(f"/api/v2/occurrences/?project_id={self.project.pk}&job={job_param}&limit=50") + + def test_returns_each_occurrence_the_job_wrote_once(self): + response = self._list(self.job.pk) + self.assertEqual(response.status_code, 200) + ids = [row["id"] for row in response.json()["results"]] + self.assertEqual(len(ids), len(set(ids))) + self.assertEqual(set(ids), {self.occ_detected.pk, self.occ_classified.pk, self.occ_multi.pk}) + self.assertEqual(response.json()["count"], 3) + + def test_other_job_matches_only_its_own_occurrences(self): + response = self._list(self.other_job.pk) + self.assertEqual(response.status_code, 200) + ids = {row["id"] for row in response.json()["results"]} + self.assertEqual(ids, {self.occ_classified.pk, self.occ_other.pk}) + + def test_non_integer_job_is_a_bad_request(self): + response = self._list("abc") + self.assertEqual(response.status_code, 400) + self.assertIn("job", response.json()) + + def test_filtered_list_query_count(self): + """Pins the query count of a filtered list over several matching rows, so a filter + that starts querying per occurrence fails here. Cachalot is off so every query counts.""" + from cachalot.api import cachalot_disabled + + disabled = cachalot_disabled() + disabled.__enter__() + try: + # Most of these are the list's existing per-row cost, not the filter: see #1461. + with self.assertNumQueries(65): + response = self._list(self.job.pk) + finally: + # cachalot_disabled() does not restore itself when the block raises. + disabled.__exit__(None, None, None) + self.assertEqual(response.status_code, 200) + self.assertEqual(response.json()["count"], 3) + + class TestCleanupNullOnlyOccurrencesCommand(TestCase): """ Covers ami/main/management/commands/cleanup_null_only_occurrences.py. From 6dbad7c3c3649b631f17a2aa38fd4a9b7bdb6294 Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Mon, 5 Oct 2026 15:17:56 -0700 Subject: [PATCH 03/12] feat: list recent processing jobs as choices for the occurrence job filter /api/v2/jobs/choices/ returns the pipeline and post-processing jobs of a project, most recently created first, in one response capped at 100, the same shape as the capture set choices. Failed jobs are included because they may have written results before failing. The capped pagination class is now shared under a generic name. Co-Authored-By: Claude Opus 5.5 (1M context) --- ami/jobs/serializers.py | 8 ++++++ ami/jobs/tests/test_jobs.py | 55 +++++++++++++++++++++++++++++++++++++ ami/jobs/views.py | 32 ++++++++++++++++++--- ami/main/api/views.py | 6 ++-- 4 files changed, 94 insertions(+), 7 deletions(-) diff --git a/ami/jobs/serializers.py b/ami/jobs/serializers.py index f53199e73..7b57047e1 100644 --- a/ami/jobs/serializers.py +++ b/ami/jobs/serializers.py @@ -192,6 +192,14 @@ class Meta: fields = ["id", "pipeline_slug"] +class JobChoiceSerializer(DefaultSerializer): + """What a job dropdown needs to name a job.""" + + class Meta: + model = Job + fields = ["id", "name", "details", "job_type_key", "created_at"] + + class MLJobTasksRequestSerializer(serializers.Serializer): """POST /jobs/{id}/tasks/ — request body sent by a processing service to fetch work. diff --git a/ami/jobs/tests/test_jobs.py b/ami/jobs/tests/test_jobs.py index 00b7934a7..01ed7431d 100644 --- a/ami/jobs/tests/test_jobs.py +++ b/ami/jobs/tests/test_jobs.py @@ -9,6 +9,7 @@ from ami.base.serializers import reverse_with_params from ami.jobs.models import ( + DataExportJob, DataStorageSyncJob, Job, JobDispatchMode, @@ -16,6 +17,7 @@ JobProgress, JobState, MLJob, + PostProcessingJob, RegroupEventsJob, SourceImageCollectionPopulateJob, ) @@ -1744,3 +1746,56 @@ def test_browsable_page_renders_number_input(self): html = response.content.decode() self.assertNotIn('