From 2c970f8371ad3992cfe71255d542372facb7fdb0 Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Tue, 6 Oct 2026 00:08:46 -0700 Subject: [PATCH 1/4] feat(ml): store vectors from a feature-only response and refuse such pipelines in ML jobs A feature-only processing service returns the detections it was sent, unchanged and still naming the detector, with an embedding attached. The new save_embedding_results() stores only those vectors on the stored detections they match by capture and box, so it can never create a detection, classification or occurrence, and it ignores the echoed detector reference. An embedding from an algorithm outside the pipeline still raises, and the per-(algorithm, key) length check still applies. Unmatched boxes are counted in the result instead of created. Pipeline gains embedding_algorithms() and is_embedding_only(). A regular ML job (and Pipeline.process_images) now refuses an embedding-only pipeline up front with a message pointing to the "Add feature vectors" admin action, instead of failing later in save_results with PipelineNotConfigured. process_detections() sends a chosen set of stored detections in one synchronous request, sharing the request building with collect_detections(). Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_0121zMVjnPsqeDFSBXRCPvMy --- ami/jobs/models.py | 1 + ami/ml/embeddings/writer.py | 88 +++++++++++++++++++++++++++++---- ami/ml/models/pipeline.py | 97 ++++++++++++++++++++++++++++++------- 3 files changed, 159 insertions(+), 27 deletions(-) diff --git a/ami/jobs/models.py b/ami/jobs/models.py index ff65f31f2..ea0f03e71 100644 --- a/ami/jobs/models.py +++ b/ami/jobs/models.py @@ -500,6 +500,7 @@ def run(cls, job: "Job"): if not job.pipeline: raise ValueError("No pipeline specified to process images in ML job") + job.pipeline.raise_if_embedding_only() job.progress.update_stage( "collect", diff --git a/ami/ml/embeddings/writer.py b/ami/ml/embeddings/writer.py index e55c72371..38321d22f 100644 --- a/ami/ml/embeddings/writer.py +++ b/ami/ml/embeddings/writer.py @@ -1,8 +1,12 @@ """Store the feature vectors a processing service returns with its detections.""" +from __future__ import annotations + import collections import contextlib +import dataclasses import logging +import typing import zlib import numpy as np @@ -12,7 +16,11 @@ from ami.ml.exceptions import PipelineNotConfigured from ami.ml.models.algorithm import Algorithm from ami.ml.models.embedding import DetectionEmbedding, as_half_precision -from ami.ml.schemas import DetectionResponse +from ami.ml.schemas import DetectionResponse, PipelineResultsResponse + +if typing.TYPE_CHECKING: + from ami.jobs.models import Job + from ami.ml.models.pipeline import Pipeline logger = logging.getLogger(__name__) @@ -56,6 +64,17 @@ def _check_embedding_dimensions( ) +@dataclasses.dataclass +class EmbeddingStoreResult: + """What one batch of returned vectors did: the rows sent to the store and how each vector fared.""" + + embeddings: list[DetectionEmbedding] = dataclasses.field(default_factory=list) + written: int = 0 + unchanged: int = 0 + unmatched: int = 0 + not_finite: int = 0 + + def create_detection_embeddings( detections: list[Detection], detection_responses: list[DetectionResponse], @@ -80,6 +99,23 @@ def create_detection_embeddings( ``EmbeddingDimensionMismatch``. ``job_id`` records the job whose results stored each vector. Returns the embeddings that were sent to the store (written or unchanged). """ + return store_detection_embeddings( + detections, detection_responses, algorithms_known, DEFAULT_EMBEDDING_KEY, logger, job_id + ).embeddings + + +def store_detection_embeddings( + detections: list[Detection], + detection_responses: list[DetectionResponse], + algorithms_known: dict[str, Algorithm], + key: str, + logger: logging.Logger = logger, + job_id: int | None = None, +) -> EmbeddingStoreResult: + """Match returned vectors to ``detections`` by image and box and store them under ``key``. + + The matching, validation and storing rules are those of ``create_detection_embeddings``. + """ by_box: dict[tuple, list[Detection]] = collections.defaultdict(list) for detection in detections: if detection.bbox is not None: @@ -113,25 +149,28 @@ def create_detection_embeddings( if not np.isfinite(as_half_precision(embedding_resp.features)).all(): not_finite += 1 continue - lengths_by_pair[(algorithm, DEFAULT_EMBEDDING_KEY)].add(len(embedding_resp.features)) - embeddings[(detection.pk, algorithm.pk, DEFAULT_EMBEDDING_KEY)] = DetectionEmbedding( + lengths_by_pair[(algorithm, key)].add(len(embedding_resp.features)) + embeddings[(detection.pk, algorithm.pk, key)] = DetectionEmbedding( detection=detection, algorithm=algorithm, - key=DEFAULT_EMBEDDING_KEY, + key=key, vector=embedding_resp.features, job_id=job_id, ) stored_by_pair = {pair: DetectionEmbedding.objects.stored_length(pair[0].pk, pair[1]) for pair in lengths_by_pair} first_writes = [pair for pair, stored in stored_by_pair.items() if stored is None] + result = EmbeddingStoreResult(embeddings=list(embeddings.values()), unmatched=unmatched, not_finite=not_finite) with contextlib.ExitStack() as stack: if first_writes: # Hold the lock until the vectors are committed, so a concurrent first writer sees them. stack.enter_context(transaction.atomic()) - for algorithm, key in first_writes: - _lock_first_write(algorithm, key) - stored_by_pair[(algorithm, key)] = DetectionEmbedding.objects.stored_length(algorithm.pk, key) + for algorithm, pair_key in first_writes: + _lock_first_write(algorithm, pair_key) + stored_by_pair[(algorithm, pair_key)] = DetectionEmbedding.objects.stored_length( + algorithm.pk, pair_key + ) _check_embedding_dimensions(lengths_by_pair, stored_by_pair) if unmatched: @@ -139,6 +178,35 @@ def create_detection_embeddings( if not_finite: logger.warning(f"Skipped {not_finite} vectors with values a half-precision vector cannot store.") if embeddings: - written, unchanged = DetectionEmbedding.objects.store(embeddings.values()) - logger.info(f"Stored {written} feature vectors ({unchanged} unchanged) for {len(detections)} detections.") - return list(embeddings.values()) + result.written, result.unchanged = DetectionEmbedding.objects.store(embeddings.values()) + logger.info( + f"Stored {result.written} feature vectors ({result.unchanged} unchanged) " + f"for {len(detections)} detections." + ) + return result + + +def save_embedding_results( + response: PipelineResultsResponse, + job: Job | None, + pipeline: Pipeline, + key: str = DEFAULT_EMBEDDING_KEY, +) -> EmbeddingStoreResult: + """Store only the vectors of a feature-only response, on detections that already exist. + + A feature-only processing service returns the detections it was sent, unchanged and still + naming the detector that made them, with an embedding attached and no classifications. Each + returned detection is matched by capture and box to a stored valid detection; nothing else is + written: no detection, classification or occurrence is created, and the echoed detector + reference is ignored. A returned box with no stored detection is counted in ``unmatched`` and + logged, never created. An embedding that names an algorithm outside ``pipeline`` raises + ``PipelineNotConfigured``, and a vector of the wrong length for its (algorithm, key) raises + ``EmbeddingDimensionMismatch``. + """ + capture_ids = {int(detection.source_image_id) for detection in response.detections} + detections = list(Detection.objects.valid().filter(source_image_id__in=capture_ids)) + algorithms_known = {algorithm.key: algorithm for algorithm in pipeline.algorithms.all()} + job_logger = job.logger if job else logger + return store_detection_embeddings( + detections, response.detections, algorithms_known, key, job_logger, job.pk if job else None + ) diff --git a/ami/ml/models/pipeline.py b/ami/ml/models/pipeline.py index b7baebd32..838fae925 100644 --- a/ami/ml/models/pipeline.py +++ b/ami/ml/models/pipeline.py @@ -40,7 +40,7 @@ ) from ami.ml.embeddings.writer import create_detection_embeddings from ami.ml.exceptions import PipelineNotConfigured -from ami.ml.models.algorithm import Algorithm, AlgorithmCategoryMap +from ami.ml.models.algorithm import Algorithm, AlgorithmCategoryMap, AlgorithmTaskType from ami.ml.schemas import ( AlgorithmConfigResponse, AlgorithmReference, @@ -411,6 +411,22 @@ def process_images( return results +def detection_request(detection: Detection, source_image_request: SourceImageRequest) -> DetectionRequest | None: + """The request form of a stored detection, or None when it has no box or no detector to name.""" + bbox = detection.get_bbox() + if not bbox or not detection.detection_algorithm: + return None + return DetectionRequest( + source_image=source_image_request, + bbox=bbox, + crop_image_url=detection.url(), + algorithm=AlgorithmReference( + name=detection.detection_algorithm.name, + key=detection.detection_algorithm.key, + ), + ) + + def collect_detections( source_image: SourceImage, source_image_request: SourceImageRequest, @@ -418,24 +434,50 @@ def collect_detections( """ Collect existing detections for a source image and send them with pipeline request. """ - detection_requests: list[DetectionRequest] = [] # Re-process all existing detections if they exist - for detection in source_image.detections.all(): - bbox = detection.get_bbox() - if bbox and detection.detection_algorithm: - detection_requests.append( - DetectionRequest( - source_image=source_image_request, - bbox=bbox, - crop_image_url=detection.url(), - algorithm=AlgorithmReference( - name=detection.detection_algorithm.name, - key=detection.detection_algorithm.key, - ), - ) - ) + built = (detection_request(detection, source_image_request) for detection in source_image.detections.all()) + return [request for request in built if request is not None] + - return detection_requests +def process_detections( + pipeline: Pipeline, + endpoint_url: str, + detections: typing.Iterable[Detection], + project_id: int, +) -> PipelineResultsResponse: + """Send stored detections, and only those, to a processing service in one synchronous request. + + Serves feature-only pipelines, which embed the detections they are sent instead of detecting. + The detections' captures (``source_image`` with its deployment and data source) and detector + (``detection_algorithm``) should be loaded with the queryset, because each is read per detection. + A failed request raises ``requests.HTTPError`` rather than returning an empty response. + """ + source_image_requests: dict[int, SourceImageRequest] = {} + detection_requests: list[DetectionRequest] = [] + for detection in detections: + source_image = detection.source_image + url = source_image.public_url() + if not url: + continue + source_image_request = source_image_requests.setdefault( + source_image.pk, SourceImageRequest(id=str(source_image.pk), url=url) + ) + request = detection_request(detection, source_image_request) + if request is not None: + detection_requests.append(request) + + request_data = PipelineRequest( + pipeline=pipeline.slug, + source_images=list(source_image_requests.values()), + config=pipeline.get_config(project_id=project_id), + detections=detection_requests, + ) + resp = create_session().post(endpoint_url, json=request_data.dict()) + if not resp.ok: + raise requests.HTTPError( + f"Failed to process {request_data.summary()}: {extract_error_message_from_response(resp)}" + ) + return PipelineResultsResponse(**resp.json()) def get_or_create_algorithm_and_category_map( @@ -1293,6 +1335,26 @@ def collect_images( reprocess_all_images=reprocess_all_images, ) + def embedding_algorithms(self) -> models.QuerySet[Algorithm]: + """The algorithms of this pipeline that produce feature vectors.""" + return self.algorithms.filter(task_type=AlgorithmTaskType.EMBEDDING.value) + + def is_embedding_only(self) -> bool: + """True when every algorithm of the pipeline produces feature vectors, and there is at least one. + + Such a pipeline detects and classifies nothing, so it can only be run through the + "Add feature vectors" task, which stores vectors on detections that already exist. + """ + algorithms = self.algorithms.all() + return bool(algorithms) and all(a.task_type == AlgorithmTaskType.EMBEDDING.value for a in algorithms) + + def raise_if_embedding_only(self) -> None: + if self.is_embedding_only(): + raise PipelineNotConfigured( + f'Pipeline "{self.name}" only produces feature vectors, so it cannot run as an ML job. ' + 'Use the "Add feature vectors" action on a capture set in the admin instead.' + ) + def choose_processing_service_for_pipeline( self, job_id: int | None, pipeline_name: str, project_id: int ) -> ProcessingService: @@ -1348,6 +1410,7 @@ def process_images( job_id: int | None = None, reprocess_all_images: bool = False, ) -> PipelineResultsResponse: + self.raise_if_embedding_only() processing_service = self.choose_processing_service_for_pipeline(job_id, self.name, project_id) if not processing_service.endpoint_url: From dbdbc33258b922d7df70c0a87c0b5178d99856e9 Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Tue, 6 Oct 2026 00:08:47 -0700 Subject: [PATCH 2/4] feat(ml): add an admin action that adds feature vectors to existing detections Admins can now select capture sets and run "Add feature vectors" with a pipeline that has an embedding algorithm. Only valid detections that lack a vector from that extractor are sent, in batches of captures, to the pipeline's processing service, and the returned vectors are stored on those detections. Captures with nothing missing are skipped, so a second run sends nothing. One job tracks the run and reports captures, detections sent, vectors stored, vectors unchanged and boxes unmatched on its stage. The task is registered like the other post-processing tasks and triggered through make_post_processing_action; no new job type, migration or UI is involved. The "still missing" reads bypass the query cache, which otherwise served the answer from before the vectors were stored. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_0121zMVjnPsqeDFSBXRCPvMy --- ami/main/admin.py | 11 + .../admin/feature_vectors_form.py | 51 +++ ami/ml/post_processing/feature_vectors.py | 112 ++++++ ami/ml/post_processing/registry.py | 2 + .../tests/test_feature_vectors.py | 358 ++++++++++++++++++ 5 files changed, 534 insertions(+) create mode 100644 ami/ml/post_processing/admin/feature_vectors_form.py create mode 100644 ami/ml/post_processing/feature_vectors.py create mode 100644 ami/ml/post_processing/tests/test_feature_vectors.py diff --git a/ami/main/admin.py b/ami/main/admin.py index 55c84701f..a96a3cc56 100644 --- a/ami/main/admin.py +++ b/ami/main/admin.py @@ -15,8 +15,10 @@ from ami.ml.models.project_pipeline_config import ProjectPipelineConfig from ami.ml.post_processing.admin.actions import make_post_processing_action from ami.ml.post_processing.admin.class_masking_form import ClassMaskingActionForm +from ami.ml.post_processing.admin.feature_vectors_form import AddFeatureVectorsActionForm from ami.ml.post_processing.admin.small_size_filter_form import SmallSizeFilterActionForm from ami.ml.post_processing.class_masking import ClassMaskingTask +from ami.ml.post_processing.feature_vectors import AddFeatureVectorsTask from ami.ml.post_processing.small_size_filter import SmallSizeFilterTask from ami.ml.tasks import remove_duplicate_classifications @@ -859,11 +861,20 @@ def populate_collection_async(self, request: HttpRequest, queryset: QuerySet[Sou f"Post-processing: {task_cls.name} on Capture Set {collection.pk}" ), ) + run_add_feature_vectors = make_post_processing_action( + AddFeatureVectorsTask, + AddFeatureVectorsActionForm, + scope_resolver=lambda collection: {"source_image_collection_id": collection.pk}, + name_resolver=lambda task_cls, collection: ( + f"Post-processing: {task_cls.name} on Capture Set {collection.pk}" + ), + ) actions = [ populate_collection, populate_collection_async, run_small_size_filter, run_class_masking, + run_add_feature_vectors, ] # Hide images many-to-many field from form. This would list all source images in the database. diff --git a/ami/ml/post_processing/admin/feature_vectors_form.py b/ami/ml/post_processing/admin/feature_vectors_form.py new file mode 100644 index 000000000..40cd11188 --- /dev/null +++ b/ami/ml/post_processing/admin/feature_vectors_form.py @@ -0,0 +1,51 @@ +from __future__ import annotations + +from django import forms + +from ami.main.models import DEFAULT_EMBEDDING_KEY +from ami.ml.models import Pipeline +from ami.ml.models.algorithm import AlgorithmTaskType +from ami.ml.post_processing.admin.forms import BasePostProcessingActionForm +from ami.ml.post_processing.feature_vectors import AddFeatureVectorsConfig + + +class AddFeatureVectorsActionForm(BasePostProcessingActionForm): + """Knobs surfaced when an admin triggers Add feature vectors. + + The pipeline choice lists only pipelines with an algorithm that produces feature vectors, + narrowed to those enabled for the projects of the selected capture sets. The valid range of + ``batch_size`` lives on ``AddFeatureVectorsConfig``; the admin action surfaces its errors inline. + """ + + pipeline_id = forms.ModelChoiceField( + queryset=Pipeline.objects.none(), + label="Pipeline", + help_text="A pipeline with a feature extractor. Only detections that lack its vectors are sent.", + ) + key = forms.CharField( + label="Vector name", + initial=DEFAULT_EMBEDDING_KEY, + help_text="The name the vectors are stored under. Keep the default unless the extractor has several outputs.", + ) + batch_size = forms.IntegerField( + label="Captures per request", + initial=AddFeatureVectorsConfig.__fields__["batch_size"].default, + help_text="How many captures, with their missing detections, go to the processing service in one request.", + ) + + def __init__(self, *args, **kwargs): + super().__init__(*args, **kwargs) + pipelines = Pipeline.objects.filter(algorithms__task_type=AlgorithmTaskType.EMBEDDING.value) + if self.scope_queryset is not None: + project_ids = set(self.scope_queryset.values_list("project_id", flat=True)) + pipelines = pipelines.filter( + project_pipeline_configs__enabled=True, project_pipeline_configs__project_id__in=project_ids + ) + self.fields["pipeline_id"].queryset = pipelines.distinct().order_by("name") + + def to_config(self) -> dict: + return { + "pipeline_id": self.cleaned_data["pipeline_id"].pk, + "key": self.cleaned_data["key"], + "batch_size": self.cleaned_data["batch_size"], + } diff --git a/ami/ml/post_processing/feature_vectors.py b/ami/ml/post_processing/feature_vectors.py new file mode 100644 index 000000000..88e069e6c --- /dev/null +++ b/ami/ml/post_processing/feature_vectors.py @@ -0,0 +1,112 @@ +import functools +import operator +import typing +from urllib.parse import urljoin + +import pydantic +from cachalot.api import cachalot_disabled + +from ami.main.models import DEFAULT_EMBEDDING_KEY, Detection, SourceImageCollection +from ami.ml.embeddings.reader import detections_missing_vectors +from ami.ml.embeddings.writer import save_embedding_results +from ami.ml.exceptions import PipelineNotConfigured +from ami.ml.models import Pipeline +from ami.ml.models.pipeline import process_detections +from ami.ml.post_processing.base import BasePostProcessingTask + + +class AddFeatureVectorsConfig(pydantic.BaseModel): + source_image_collection_id: int + pipeline_id: int + key: str = DEFAULT_EMBEDDING_KEY + # Captures per request. Each request carries every missing detection of its captures. + batch_size: int = pydantic.Field(default=10, ge=1, le=100) + + class Config: + extra = "forbid" + + +class AddFeatureVectorsTask(BasePostProcessingTask): + """Add feature vectors to detections that already exist, without detecting or classifying again. + + For the captures of one capture set, sends only the valid detections that lack a vector from + the pipeline's feature extractor to the processing service, and stores the returned vectors + on those detections. Running it twice sends nothing the second time. + """ + + key = "add_feature_vectors" + name = "Add feature vectors" + config_schema = AddFeatureVectorsConfig + + def _missing_detections(self, collection: SourceImageCollection, pipeline: Pipeline, key: str): + """Valid detections of the capture set that lack a vector from any of the pipeline's extractors.""" + detections = Detection.objects.valid().filter(source_image__collections=collection) + return functools.reduce( + operator.or_, + ( + detections_missing_vectors(detections, extractor.pk, key) + for extractor in pipeline.embedding_algorithms() + ), + ) + + def run(self) -> None: + config = typing.cast(AddFeatureVectorsConfig, self.config) + try: + collection = SourceImageCollection.objects.get(pk=config.source_image_collection_id) + pipeline = Pipeline.objects.get(pk=config.pipeline_id) + except (SourceImageCollection.DoesNotExist, Pipeline.DoesNotExist) as err: + self.logger.error(str(err)) + raise ValueError(str(err)) from err + if not pipeline.embedding_algorithms().exists(): + msg = f'Pipeline "{pipeline.name}" has no algorithm that produces feature vectors (task type "embedding").' + self.logger.error(msg) + raise ValueError(msg) + + missing = self._missing_detections(collection, pipeline, config.key) + # The "what is still missing" reads bypass the query cache: it can keep serving the answer from + # before the vectors were stored, which would make a second run send everything again. + with cachalot_disabled(): + capture_ids = sorted(missing.order_by().values_list("source_image_id", flat=True).distinct()) + self.logger.info( + f"=== Starting {self.name}: {len(capture_ids)} captures of capture set {collection.pk} " + f"have detections without a vector from pipeline {pipeline} ===" + ) + + if not capture_ids: + self.logger.info(f"=== Completed {self.name}: nothing to add ===") + return + processing_service = pipeline.choose_processing_service_for_pipeline( + self.job.pk if self.job else None, pipeline.name, collection.project_id + ) + if not processing_service.endpoint_url: + raise PipelineNotConfigured( + f"No endpoint URL configured for this pipeline's processing service ({processing_service})" + ) + endpoint_url = urljoin(processing_service.endpoint_url, "/process") + + totals = { + "Captures": 0, + "Detections sent": 0, + "Vectors stored": 0, + "Vectors unchanged": 0, + "Boxes unmatched": 0, + } + self.report_stage_metrics(totals) + for start in range(0, len(capture_ids), config.batch_size): + with cachalot_disabled(): + batch = list( + missing.filter(source_image_id__in=capture_ids[start : start + config.batch_size]) + .select_related("source_image__deployment__data_source", "detection_algorithm") + .order_by("source_image_id", "pk") + ) + response = process_detections(pipeline, endpoint_url, batch, collection.project_id) + result = save_embedding_results(response, self.job, pipeline, config.key) + totals["Captures"] += len({detection.source_image_id for detection in batch}) + totals["Detections sent"] += len(batch) + totals["Vectors stored"] += result.written + totals["Vectors unchanged"] += result.unchanged + totals["Boxes unmatched"] += result.unmatched + self.update_progress(min(start + config.batch_size, len(capture_ids)) / len(capture_ids)) + self.report_stage_metrics(totals) + + self.logger.info(f"=== Completed {self.name}: {totals} ===") diff --git a/ami/ml/post_processing/registry.py b/ami/ml/post_processing/registry.py index 308be18ae..c4a16db33 100644 --- a/ami/ml/post_processing/registry.py +++ b/ami/ml/post_processing/registry.py @@ -1,10 +1,12 @@ # Registry of available post-processing tasks from ami.ml.post_processing.class_masking import ClassMaskingTask +from ami.ml.post_processing.feature_vectors import AddFeatureVectorsTask from ami.ml.post_processing.small_size_filter import SmallSizeFilterTask POSTPROCESSING_TASKS = { SmallSizeFilterTask.key: SmallSizeFilterTask, ClassMaskingTask.key: ClassMaskingTask, + AddFeatureVectorsTask.key: AddFeatureVectorsTask, } diff --git a/ami/ml/post_processing/tests/test_feature_vectors.py b/ami/ml/post_processing/tests/test_feature_vectors.py new file mode 100644 index 000000000..7f69bb52f --- /dev/null +++ b/ami/ml/post_processing/tests/test_feature_vectors.py @@ -0,0 +1,358 @@ +"""Domain tests for the "Add feature vectors" post-processing task. + +The task sends detections that already exist, and only those that lack a vector, to a feature-only +processing service and stores the vectors that come back. The service is stubbed with a response in +the shape ADC's feature-only pipeline returns: each requested detection echoed unchanged, still +naming the detector, with an embedding attached and no classifications. +""" +import datetime +from unittest import mock + +from cachalot.api import cachalot_disabled +from django.contrib import admin as django_admin +from django.db import connection +from django.test import Client, TestCase +from django.test.utils import CaptureQueriesContext +from django.urls import reverse + +from ami.jobs.models import Job +from ami.main.models import Classification, Detection, Occurrence, SourceImage, SourceImageCollection +from ami.ml.exceptions import PipelineNotConfigured +from ami.ml.models import Algorithm, DetectionEmbedding, Pipeline, ProcessingService +from ami.ml.models.algorithm import AlgorithmTaskType +from ami.ml.models.pipeline import get_or_create_algorithm_and_category_map +from ami.ml.post_processing.feature_vectors import AddFeatureVectorsTask +from ami.tests.fixtures.main import setup_test_project +from ami.tests.fixtures.ml import ALGORITHM_CHOICES +from ami.users.models import User + +DETECTOR = ALGORITHM_CHOICES["random-detector"] +LENGTH = 8 +VECTOR = [0.5] * LENGTH + + +class FakeFeaturePipelineService: + """Stands in for the processing service's synchronous ``/process`` endpoint. + + Mirrors ``run_feature_pipeline``: the requested detections come back with their box and detector + reference unchanged, an ``embeddings`` entry from the extractor, and no classifications. + """ + + def __init__(self, extractor_key: str): + self.extractor_key = extractor_key + self.requests: list[dict] = [] + self.extra_detections: list[dict] = [] + self.embedding_key_override: str | None = None + self.status_ok = True + + def post(self, url: str, json: dict): + self.requests.append(json) + key = self.embedding_key_override or self.extractor_key + detections = [ + { + "source_image_id": requested["source_image"]["id"], + "bbox": requested["bbox"], + "algorithm": requested["algorithm"], + "crop_image_url": requested["crop_image_url"], + "timestamp": datetime.datetime.now().isoformat(), + "classifications": [], + "embeddings": [{"algorithm": {"name": "Extractor", "key": key}, "features": VECTOR}], + } + for requested in json["detections"] + ] + self.extra_detections + body = { + "pipeline": json["pipeline"], + "total_time": 0.1, + "source_images": [{"id": image["id"], "url": image["url"]} for image in json["source_images"]], + "detections": detections, + } + return mock.Mock(ok=self.status_ok, json=lambda: body, status_code=200 if self.status_ok else 500) + + @property + def detections_sent(self) -> int: + return sum(len(request["detections"]) for request in self.requests) + + +class FeatureVectorsFixture: + """A project with a feature-only pipeline, its processing service stubbed, and captures with detections.""" + + def _set_up(self) -> None: + self.project, self.deployment = setup_test_project(reuse=False) + self.detector = get_or_create_algorithm_and_category_map(DETECTOR) + self.extractor = Algorithm.objects.create( + name="Extractor", key="extractor", task_type=AlgorithmTaskType.EMBEDDING.value + ) + self.pipeline = Pipeline.objects.create(name="Feature pipeline") + self.pipeline.algorithms.set([self.extractor]) + # Creating a service checks its status over the network right away; there is no network here. + status_patcher = mock.patch.object(ProcessingService, "get_status") + status_patcher.start() + self.addCleanup(status_patcher.stop) # type: ignore[attr-defined] + self.service = ProcessingService.objects.create( + name="Feature service", endpoint_url="http://features.test:2000", last_seen_live=True + ) + ProcessingService.objects.filter(pk=self.service.pk).update(last_seen_live=True) + self.service.projects.add(self.project) + self.service.pipelines.add(self.pipeline) + self.fake = FakeFeaturePipelineService(self.extractor.key) + patcher = mock.patch("ami.ml.models.pipeline.create_session", return_value=self.fake) + patcher.start() + self.addCleanup(patcher.stop) # type: ignore[attr-defined] + self.captures = 0 + + def _collection(self, captures: int, detections_per_capture: int = 2) -> SourceImageCollection: + images = [] + for _ in range(captures): + self.captures += 1 + image = SourceImage.objects.create( + path=f"fv-{self.captures}-20240101{self.captures:06d}.jpg", + deployment=self.deployment, + project=self.project, + public_base_url="http://images.test/", + ) + for i in range(detections_per_capture): + Detection.objects.create( + source_image=image, + bbox=[20.0 * i, 0.0, 20.0 * i + 10, 10.0], + detection_algorithm=self.detector, + ) + images.append(image) + collection = SourceImageCollection.objects.create( + name=f"fv collection {self.captures}", project=self.project, method="manual" + ) + collection.images.set(images) + return collection + + def _config(self, collection: SourceImageCollection, **extra) -> dict: + return {"source_image_collection_id": collection.pk, "pipeline_id": self.pipeline.pk, **extra} + + def _job(self, collection: SourceImageCollection, **extra) -> Job: + return Job.objects.create( + name="Add feature vectors", + project=self.project, + job_type_key="post_processing", + params={"task": AddFeatureVectorsTask.key, "config": self._config(collection, **extra)}, + ) + + def _run(self, collection: SourceImageCollection, **extra) -> Job: + job = self._job(collection, **extra) + job.run() + job.refresh_from_db() + return job + + def _metrics(self, job: Job) -> dict[str, int]: + return {param.name: param.value for param in job.progress.stages[0].params} + + +class TestAddFeatureVectorsTask(FeatureVectorsFixture, TestCase): + def setUp(self) -> None: + self._set_up() + + def test_only_detections_missing_a_vector_are_sent(self): + collection = self._collection(captures=1, detections_per_capture=3) + stored, *missing = list(Detection.objects.filter(source_image__collections=collection).order_by("pk")) + DetectionEmbedding.objects.create(detection=stored, algorithm=self.extractor, key="embedding", vector=VECTOR) + + self._run(collection) + + sent = [d["bbox"] for request in self.fake.requests for d in request["detections"]] + self.assertCountEqual( + sent, [{"x1": m.bbox[0], "y1": m.bbox[1], "x2": m.bbox[2], "y2": m.bbox[3]} for m in missing] + ) + self.assertEqual(DetectionEmbedding.objects.filter(algorithm=self.extractor).count(), 3) + + def test_captures_with_nothing_missing_are_skipped(self): + collection = self._collection(captures=2) + done, todo = list(collection.images.order_by("pk")) + for detection in done.detections.all(): + DetectionEmbedding.objects.create( + detection=detection, algorithm=self.extractor, key="embedding", vector=VECTOR + ) + + job = self._run(collection) + + (request,) = self.fake.requests + self.assertEqual([image["id"] for image in request["source_images"]], [str(todo.pk)]) + self.assertEqual(self._metrics(job)["Captures"], 1) + + def test_vectors_are_stored_on_the_detections_they_were_sent_for_and_reported(self): + collection = self._collection(captures=3) + + job = self._run(collection) + + detections = Detection.objects.filter(source_image__collections=collection) + self.assertEqual( + set(DetectionEmbedding.objects.filter(algorithm=self.extractor).values_list("detection_id", flat=True)), + set(detections.values_list("pk", flat=True)), + ) + self.assertEqual( + self._metrics(job), + { + "Captures": 3, + "Detections sent": 6, + "Vectors stored": 6, + "Vectors unchanged": 0, + "Boxes unmatched": 0, + }, + ) + + def test_no_detections_classifications_or_occurrences_are_created(self): + collection = self._collection(captures=2) + before = (Detection.objects.count(), Classification.objects.count(), Occurrence.objects.count()) + + self._run(collection) + + self.assertEqual( + (Detection.objects.count(), Classification.objects.count(), Occurrence.objects.count()), before + ) + + def test_a_returned_box_with_no_stored_detection_is_counted_and_not_created(self): + collection = self._collection(captures=1) + image = collection.images.get() + self.fake.extra_detections = [ + { + "source_image_id": str(image.pk), + "bbox": {"x1": 500.0, "y1": 500.0, "x2": 510.0, "y2": 510.0}, + "algorithm": {"name": DETECTOR.name, "key": DETECTOR.key}, + "timestamp": datetime.datetime.now().isoformat(), + "embeddings": [{"algorithm": {"name": "Extractor", "key": self.extractor.key}, "features": VECTOR}], + } + ] + + job = self._run(collection) + + self.assertEqual(self._metrics(job)["Boxes unmatched"], 1) + self.assertEqual(Detection.objects.filter(source_image=image).count(), 2) + self.assertEqual(DetectionEmbedding.objects.count(), 2) + + def test_an_embedding_from_an_algorithm_outside_the_pipeline_raises_and_stores_nothing(self): + collection = self._collection(captures=1) + self.fake.embedding_key_override = "not-in-pipeline" + + with self.assertRaises(PipelineNotConfigured): + self._run(collection) + + self.assertEqual(DetectionEmbedding.objects.count(), 0) + + def test_a_second_run_sends_nothing(self): + collection = self._collection(captures=2) + self._run(collection) + requests_after_first = len(self.fake.requests) + + self._run(collection) + + self.assertEqual(len(self.fake.requests), requests_after_first) + self.assertEqual(DetectionEmbedding.objects.count(), 4) + + def test_a_pipeline_without_an_embedding_algorithm_is_refused_before_any_request(self): + collection = self._collection(captures=1) + self.pipeline.algorithms.set([self.detector]) + + with self.assertRaisesMessage(ValueError, "produces feature vectors"): + self._run(collection) + + self.assertEqual(self.fake.requests, []) + + def test_a_failed_request_raises_instead_of_reporting_success(self): + collection = self._collection(captures=1) + self.fake.status_ok = False + with mock.patch("ami.ml.models.pipeline.extract_error_message_from_response", return_value="boom"): + with self.assertRaisesMessage(Exception, "boom"): + self._run(collection) + self.assertEqual(DetectionEmbedding.objects.count(), 0) + + def test_requests_are_made_per_batch_of_captures(self): + collection = self._collection(captures=5) + self._run(collection, batch_size=2) + self.assertEqual([len(r["source_images"]) for r in self.fake.requests], [2, 2, 1]) + + def test_queries_do_not_grow_with_the_number_of_captures(self): + """A batch of 2 captures and a batch of 8 take the same number of queries.""" + self._run(self._collection(captures=1)) # first run creates the task's algorithm row and similar one-offs + small, large = self._collection(captures=2), self._collection(captures=8) + # Cold counts: a query the cache served for the first run would otherwise be missing from the second. + with cachalot_disabled(): + with CaptureQueriesContext(connection) as small_queries: + self._run(small, batch_size=100) + with CaptureQueriesContext(connection) as large_queries: + self._run(large, batch_size=100) + self.assertEqual(len(large_queries), len(small_queries)) + + +class TestEmbeddingOnlyPipelineGuard(FeatureVectorsFixture, TestCase): + def setUp(self) -> None: + self._set_up() + + def test_pipeline_knows_when_it_only_produces_vectors(self): + self.assertTrue(self.pipeline.is_embedding_only()) + self.pipeline.algorithms.add(self.detector) + self.assertFalse(self.pipeline.is_embedding_only()) + self.assertFalse(Pipeline.objects.create(name="Empty").is_embedding_only()) + + def test_a_regular_ml_job_with_an_embedding_only_pipeline_fails_early_pointing_to_the_task(self): + collection = self._collection(captures=1) + job = Job.objects.create( + name="ML job", + project=self.project, + job_type_key="ml", + pipeline=self.pipeline, + source_image_collection=collection, + ) + with self.assertRaisesMessage(PipelineNotConfigured, "Add feature vectors"): + job.run() + self.assertEqual(self.fake.requests, []) + + +class TestAddFeatureVectorsAdmin(FeatureVectorsFixture, TestCase): + def setUp(self) -> None: + self._set_up() + self.superuser = User.objects.create_superuser(email="afv-admin@example.com", password="x") + self.client = Client() + self.client.force_login(self.superuser) + from ami.ml.models import ProjectPipelineConfig + + ProjectPipelineConfig.objects.create(project=self.project, pipeline=self.pipeline, enabled=True) + self.other_pipeline = Pipeline.objects.create(name="Classifier only pipeline") + self.other_pipeline.algorithms.set([self.detector]) + + def _post(self, collections: list[SourceImageCollection], data: dict): + return self.client.post( + reverse("admin:main_sourceimagecollection_changelist"), + { + "action": "run_add_feature_vectors", + django_admin.helpers.ACTION_CHECKBOX_NAME: [str(c.pk) for c in collections], + **data, + }, + ) + + def test_the_form_lists_only_pipelines_with_an_embedding_algorithm(self): + response = self._post([self._collection(captures=1)], {}) + self.assertEqual(response.status_code, 200) + self.assertContains(response, "Run Add feature vectors") + self.assertContains(response, "Feature pipeline") + self.assertNotContains(response, "Classifier only pipeline") + + def test_it_enqueues_one_job_per_capture_set_with_the_config_on_the_job(self): + first, second = self._collection(captures=1), self._collection(captures=1) + + response = self._post( + [first, second], {"confirm": "1", "pipeline_id": self.pipeline.pk, "key": "embedding", "batch_size": 5} + ) + + self.assertEqual(response.status_code, 302) + jobs = Job.objects.filter(job_type_key="post_processing", params__task="add_feature_vectors") + self.assertEqual(jobs.count(), 2) + self.assertCountEqual( + [job.params["config"]["source_image_collection_id"] for job in jobs], [first.pk, second.pk] + ) + for job in jobs: + self.assertEqual(job.params["config"]["pipeline_id"], self.pipeline.pk) + self.assertEqual(job.params["config"]["batch_size"], 5) + + def test_an_out_of_range_batch_size_is_shown_on_the_form_and_creates_no_job(self): + response = self._post( + [self._collection(captures=1)], + {"confirm": "1", "pipeline_id": self.pipeline.pk, "key": "embedding", "batch_size": 0}, + ) + self.assertEqual(response.status_code, 200) + self.assertEqual(Job.objects.filter(params__task="add_feature_vectors").count(), 0) From a04626e6c4563dd4b8158b35929e94fee98b4e02 Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Tue, 6 Oct 2026 00:25:16 -0700 Subject: [PATCH 3/4] fix(ml): read what is still missing through the query cache again The reader that finds detections without a vector now invalidates correctly when vectors are stored, so the add-feature-vectors task no longer needs to bypass the query cache for those reads. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_0121zMVjnPsqeDFSBXRCPvMy --- ami/ml/post_processing/feature_vectors.py | 17 ++++++----------- 1 file changed, 6 insertions(+), 11 deletions(-) diff --git a/ami/ml/post_processing/feature_vectors.py b/ami/ml/post_processing/feature_vectors.py index 88e069e6c..ec3c82aa7 100644 --- a/ami/ml/post_processing/feature_vectors.py +++ b/ami/ml/post_processing/feature_vectors.py @@ -4,7 +4,6 @@ from urllib.parse import urljoin import pydantic -from cachalot.api import cachalot_disabled from ami.main.models import DEFAULT_EMBEDDING_KEY, Detection, SourceImageCollection from ami.ml.embeddings.reader import detections_missing_vectors @@ -63,10 +62,7 @@ def run(self) -> None: raise ValueError(msg) missing = self._missing_detections(collection, pipeline, config.key) - # The "what is still missing" reads bypass the query cache: it can keep serving the answer from - # before the vectors were stored, which would make a second run send everything again. - with cachalot_disabled(): - capture_ids = sorted(missing.order_by().values_list("source_image_id", flat=True).distinct()) + capture_ids = sorted(missing.order_by().values_list("source_image_id", flat=True).distinct()) self.logger.info( f"=== Starting {self.name}: {len(capture_ids)} captures of capture set {collection.pk} " f"have detections without a vector from pipeline {pipeline} ===" @@ -93,12 +89,11 @@ def run(self) -> None: } self.report_stage_metrics(totals) for start in range(0, len(capture_ids), config.batch_size): - with cachalot_disabled(): - batch = list( - missing.filter(source_image_id__in=capture_ids[start : start + config.batch_size]) - .select_related("source_image__deployment__data_source", "detection_algorithm") - .order_by("source_image_id", "pk") - ) + batch = list( + missing.filter(source_image_id__in=capture_ids[start : start + config.batch_size]) + .select_related("source_image__deployment__data_source", "detection_algorithm") + .order_by("source_image_id", "pk") + ) response = process_detections(pipeline, endpoint_url, batch, collection.project_id) result = save_embedding_results(response, self.job, pipeline, config.key) totals["Captures"] += len({detection.source_image_id for detection in batch}) From d93e2b6f08f7d397777a4073f611f9824347561a Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Tue, 6 Oct 2026 16:02:04 -0700 Subject: [PATCH 4/4] test(ml): build the feature-vector task tests' fixtures once per class The three test classes for the add-feature-vectors task rebuilt their project, pipeline and service in setUp for every test. They now build the data once in setUpTestData; only the fake processing-service session and the status-check stub stay per test. Measured on the 16 tests in test_feature_vectors.py: 13.5 s before, 3.6 s after. The test count and results are unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_0121zMVjnPsqeDFSBXRCPvMy --- .../tests/test_feature_vectors.py | 67 ++++++++++++------- 1 file changed, 41 insertions(+), 26 deletions(-) diff --git a/ami/ml/post_processing/tests/test_feature_vectors.py b/ami/ml/post_processing/tests/test_feature_vectors.py index 7f69bb52f..150ccdcbc 100644 --- a/ami/ml/post_processing/tests/test_feature_vectors.py +++ b/ami/ml/post_processing/tests/test_feature_vectors.py @@ -18,11 +18,11 @@ from ami.jobs.models import Job from ami.main.models import Classification, Detection, Occurrence, SourceImage, SourceImageCollection from ami.ml.exceptions import PipelineNotConfigured -from ami.ml.models import Algorithm, DetectionEmbedding, Pipeline, ProcessingService +from ami.ml.models import Algorithm, DetectionEmbedding, Pipeline, ProcessingService, ProjectPipelineConfig from ami.ml.models.algorithm import AlgorithmTaskType from ami.ml.models.pipeline import get_or_create_algorithm_and_category_map from ami.ml.post_processing.feature_vectors import AddFeatureVectorsTask -from ami.tests.fixtures.main import setup_test_project +from ami.tests.fixtures.main import no_processing_service_http, setup_test_project from ami.tests.fixtures.ml import ALGORITHM_CHOICES from ami.users.models import User @@ -76,29 +76,34 @@ def detections_sent(self) -> int: class FeatureVectorsFixture: """A project with a feature-only pipeline, its processing service stubbed, and captures with detections.""" - def _set_up(self) -> None: - self.project, self.deployment = setup_test_project(reuse=False) - self.detector = get_or_create_algorithm_and_category_map(DETECTOR) - self.extractor = Algorithm.objects.create( - name="Extractor", key="extractor", task_type=AlgorithmTaskType.EMBEDDING.value - ) - self.pipeline = Pipeline.objects.create(name="Feature pipeline") - self.pipeline.algorithms.set([self.extractor]) + @classmethod + def _set_up_data(cls) -> None: # Creating a service checks its status over the network right away; there is no network here. + with no_processing_service_http(): + cls.project, cls.deployment = setup_test_project(reuse=False) + cls.detector = get_or_create_algorithm_and_category_map(DETECTOR) + cls.extractor = Algorithm.objects.create( + name="Extractor", key="extractor", task_type=AlgorithmTaskType.EMBEDDING.value + ) + cls.pipeline = Pipeline.objects.create(name="Feature pipeline") + cls.pipeline.algorithms.set([cls.extractor]) + cls.service = ProcessingService.objects.create( + name="Feature service", endpoint_url="http://features.test:2000", last_seen_live=True + ) + ProcessingService.objects.filter(pk=cls.service.pk).update(last_seen_live=True) + cls.service.projects.add(cls.project) + cls.service.pipelines.add(cls.pipeline) + cls.captures = 0 + + def _set_up_stubs(self) -> None: + """Per-test stand-ins for the processing service: a fake session and no status checks.""" status_patcher = mock.patch.object(ProcessingService, "get_status") status_patcher.start() self.addCleanup(status_patcher.stop) # type: ignore[attr-defined] - self.service = ProcessingService.objects.create( - name="Feature service", endpoint_url="http://features.test:2000", last_seen_live=True - ) - ProcessingService.objects.filter(pk=self.service.pk).update(last_seen_live=True) - self.service.projects.add(self.project) - self.service.pipelines.add(self.pipeline) self.fake = FakeFeaturePipelineService(self.extractor.key) patcher = mock.patch("ami.ml.models.pipeline.create_session", return_value=self.fake) patcher.start() self.addCleanup(patcher.stop) # type: ignore[attr-defined] - self.captures = 0 def _collection(self, captures: int, detections_per_capture: int = 2) -> SourceImageCollection: images = [] @@ -145,8 +150,12 @@ def _metrics(self, job: Job) -> dict[str, int]: class TestAddFeatureVectorsTask(FeatureVectorsFixture, TestCase): + @classmethod + def setUpTestData(cls) -> None: + cls._set_up_data() + def setUp(self) -> None: - self._set_up() + self._set_up_stubs() def test_only_detections_missing_a_vector_are_sent(self): collection = self._collection(captures=1, detections_per_capture=3) @@ -280,8 +289,12 @@ def test_queries_do_not_grow_with_the_number_of_captures(self): class TestEmbeddingOnlyPipelineGuard(FeatureVectorsFixture, TestCase): + @classmethod + def setUpTestData(cls) -> None: + cls._set_up_data() + def setUp(self) -> None: - self._set_up() + self._set_up_stubs() def test_pipeline_knows_when_it_only_produces_vectors(self): self.assertTrue(self.pipeline.is_embedding_only()) @@ -304,16 +317,18 @@ def test_a_regular_ml_job_with_an_embedding_only_pipeline_fails_early_pointing_t class TestAddFeatureVectorsAdmin(FeatureVectorsFixture, TestCase): + @classmethod + def setUpTestData(cls) -> None: + cls._set_up_data() + cls.superuser = User.objects.create_superuser(email="afv-admin@example.com", password="x") + ProjectPipelineConfig.objects.create(project=cls.project, pipeline=cls.pipeline, enabled=True) + cls.other_pipeline = Pipeline.objects.create(name="Classifier only pipeline") + cls.other_pipeline.algorithms.set([cls.detector]) + def setUp(self) -> None: - self._set_up() - self.superuser = User.objects.create_superuser(email="afv-admin@example.com", password="x") + self._set_up_stubs() self.client = Client() self.client.force_login(self.superuser) - from ami.ml.models import ProjectPipelineConfig - - ProjectPipelineConfig.objects.create(project=self.project, pipeline=self.pipeline, enabled=True) - self.other_pipeline = Pipeline.objects.create(name="Classifier only pipeline") - self.other_pipeline.algorithms.set([self.detector]) def _post(self, collections: list[SourceImageCollection], data: dict): return self.client.post(