Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
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
1 change: 1 addition & 0 deletions ami/jobs/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
11 changes: 11 additions & 0 deletions ami/main/admin.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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.
Expand Down
88 changes: 78 additions & 10 deletions ami/ml/embeddings/writer.py
Original file line number Diff line number Diff line change
@@ -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
Expand All @@ -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__)

Expand Down Expand Up @@ -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],
Expand All @@ -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:
Expand Down Expand Up @@ -113,32 +149,64 @@ 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:
logger.warning(f"Skipped the vectors of {unmatched} returned boxes that match no single stored detection.")
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
)
97 changes: 80 additions & 17 deletions ami/ml/models/pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -411,31 +411,73 @@ 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,
) -> list[DetectionRequest]:
"""
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(
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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:
Expand Down
51 changes: 51 additions & 0 deletions ami/ml/post_processing/admin/feature_vectors_form.py
Original file line number Diff line number Diff line change
@@ -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"],
}
Loading