Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
49 commits
Select commit Hold shift + click to select a range
eaa654a
feat(ml): store a feature vector for every detection the processing s…
mihow Sep 15, 2026
95e2bdb
feat(tracking): compare detections by their stored embeddings in trac…
mihow Sep 28, 2026
33f4144
fix(exports): count a detection embedding as a feature vector in the …
mihow Sep 23, 2026
4eae25e
feat(embeddings): name the stored vector "vector" and record the job …
mihow Sep 23, 2026
e33406b
test: move the capture feature-count test back into its own test case
mihow Sep 23, 2026
91dbe9a
feat(occurrences): store an occurrence's algorithm results and review…
mihow Sep 23, 2026
4cd6942
feat(occurrences): record tracking, class masking, the size filter an…
mihow Sep 23, 2026
dea7a5e
feat(occurrences): show everything that happened to an occurrence at …
mihow Sep 23, 2026
1350b81
fix(post-processing): save the size filter's last batch when its fina…
mihow Sep 23, 2026
f0d1f57
fix(occurrences): open the history of an occurrence the default filte…
mihow Sep 23, 2026
dc5345c
fix(occurrences): keep track reviews true to who confirmed which dete…
mihow Sep 23, 2026
655a7c9
feat(ui): read an occurrence's history as one list of timeline cards
mihow Sep 23, 2026
f01bfdf
feat(ui): show identifications, predictions, algorithm results and re…
mihow Sep 23, 2026
dd92ce5
fix(ui): keep a masked classifier's prediction agreeable when the det…
mihow Sep 23, 2026
94a6843
fix(ui): stop a prediction without an algorithm from linking to a bla…
mihow Sep 23, 2026
d8c2b13
fix(occurrences): let the server say whether a track changed since it…
mihow Sep 23, 2026
6f2b642
fix(ui): advance to the next occurrence after confirming from a timel…
mihow Sep 23, 2026
cf02f55
fix(ui): show the occurrence timeline without a spinner while its his…
mihow Sep 23, 2026
8d672f2
fix(ui): stop labelling signed-in users without a display name as ano…
mihow Sep 23, 2026
370ab68
fix(ui): translate the determination change on algorithm result cards
mihow Sep 23, 2026
a3e31dc
fix(occurrences): post a track review only when the detections or the…
mihow Sep 23, 2026
b52eb97
fix(occurrences): show one prediction per algorithm in the occurrence…
mihow Sep 23, 2026
271515c
fix(ui): advance to the next occurrence only when a timeline card con…
mihow Sep 23, 2026
1e35f58
fix(ui): label unnamed users the same way in the session viewer and p…
mihow Sep 23, 2026
8c89244
test: count the occurrence history queries in the multi-frame merge b…
mihow Sep 23, 2026
847b839
fix(ui): show one prediction per algorithm in the timeline before its…
mihow Sep 23, 2026
6f14d0c
style: collapse the tracking task import in the capture matcher onto …
mihow Sep 28, 2026
e2fc806
test: read stored embedding vectors as plain lists
mihow Sep 28, 2026
f8b3eb8
feat(embeddings): store vectors of any length and fill them in for ex…
mihow Sep 28, 2026
0f644ae
feat(tracking): compare the project's default feature extractor when …
mihow Sep 28, 2026
2b80254
docs: describe how detection embeddings are stored, filled in and cho…
mihow Sep 28, 2026
a190dd7
docs: index the detection embeddings reference [skip ci]
mihow Sep 28, 2026
d86a402
fix(embeddings): drop the fixed vector length on databases that ran a…
mihow Sep 28, 2026
001485c
fix(tracking): default to the feature extractor that covers the most …
mihow Sep 28, 2026
fac7871
fix(ml): keep the Collect stage heartbeat when a feature-only pipelin…
mihow Sep 28, 2026
9831afe
fix(ml): never queue a feature-only task without boxes to embed
mihow Sep 28, 2026
dee996e
docs: describe the coverage-based default extractor, the vector-lengt…
mihow Sep 28, 2026
91b9f62
fix(ml): treat a pipeline as feature-only when it lists a detector be…
mihow Sep 28, 2026
a122b89
fix(ml): stop reprocessing classified captures when a classifier pipe…
mihow Sep 28, 2026
9b2ae7c
test: accept embedding vectors of any length in the results schema
mihow Sep 28, 2026
23b44ca
fix(ml): only treat a pipeline as feature-only when every other algor…
mihow Sep 29, 2026
ccf3f3c
Merge #1432 into the history and embeddings branch
mihow Sep 29, 2026
f5aae46
fix(ml): read embedding vectors only from the "features" key
mihow Sep 29, 2026
5b95f02
fix(ml): fail a feature-only job whose returned boxes match no detection
mihow Sep 29, 2026
32da6cf
fix(ml): fail a feature-only batch that stores no vector, and count d…
mihow Sep 29, 2026
9cbb239
Merge the tracking UI branch (with main) into feat/occurrence-history…
mihow Sep 30, 2026
5a7f2cb
Merge the tracking UI branch into feat/occurrence-history-and-embeddings
mihow Sep 30, 2026
c84518b
Merge the tracking UI branch into feat/occurrence-history-and-embeddings
mihow Sep 30, 2026
fd559de
Merge the tracking server PR's review fixes, through the tracking UI PR
mihow Oct 2, 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
14 changes: 14 additions & 0 deletions ami/exports/tests.py
Original file line number Diff line number Diff line change
Expand Up @@ -688,6 +688,20 @@ def test_feature_vector_and_next_detection(self):
if pgvector_is_available():
self.assertEqual(rows[detections[0].pk]["has_feature_vector"], "true")

def test_an_embedding_alone_counts_as_a_feature_vector(self):
from ami.main.models import DetectionEmbedding
from ami.ml.models import Algorithm
from ami.tests.fixtures.tracking import pgvector_is_available

if not pgvector_is_available():
self.skipTest("This database cannot store embeddings.")
detection = self.occurrences[0].detections.order_by("source_image__timestamp").last()
algorithm = Algorithm.objects.create(name="Embedding model", key="embedding-model")
DetectionEmbedding.objects.create(detection=detection, algorithm=algorithm, vector=[0.1] * 2048)

rows = {int(row["detection_id"]): row for row in self._rows()}
self.assertEqual(rows[detection.pk]["has_feature_vector"], "true")

def test_query_count_is_one_pair_per_chunk(self):
from django.db import connection
from django.test.utils import CaptureQueriesContext
Expand Down
12 changes: 7 additions & 5 deletions ami/exports/tracks.py
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,9 @@
from collections.abc import Callable, Iterator

from django.db import models
from django.db.models import Exists, OuterRef, Subquery
from django.db.models import Exists, ExpressionWrapper, OuterRef, Subquery

from ami.main.models import BEST_MACHINE_PREDICTION_ORDER, Classification, Detection, Occurrence
from ami.main.models import BEST_MACHINE_PREDICTION_ORDER, Classification, Detection, DetectionEmbedding, Occurrence
from ami.main.models_future.tracks import CAPTURE_ORDER

TRACKS_CSV_COLUMNS: typing.Final = (
Expand Down Expand Up @@ -67,9 +67,11 @@ def _detections_for(occurrence_ids: list[int]) -> models.QuerySet:
.annotate(
label=Subquery(best_classification.values("taxon__name")[:1]),
label_score=Subquery(best_classification.values("score")[:1]),
# Same notion as ClassificationQuerySet.with_has_features(), per detection.
has_feature_vector=Exists(
Classification.objects.filter(detection=OuterRef("pk"), features_2048__isnull=False)
# Same notion as DetectionQuerySet.has_vector(): an embedding or a classification vector.
has_feature_vector=ExpressionWrapper(
Exists(DetectionEmbedding.objects.filter(detection=OuterRef("pk")))
| Exists(Classification.objects.filter(detection=OuterRef("pk"), features_2048__isnull=False)),
output_field=models.BooleanField(),
),
)
.order_by("occurrence_id", *CAPTURE_ORDER)
Expand Down
40 changes: 39 additions & 1 deletion ami/jobs/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
from redis.exceptions import RedisError

from ami.main.checks.schemas import IntegrityCheckResult
from ami.ml.exceptions import FeatureResultsStoredNothing
from ami.ml.orchestration.async_job_state import AsyncJobStateManager
from ami.ml.orchestration.nats_queue import ConsumerState, TaskQueueManager
from ami.ml.schemas import PipelineResultsError, PipelineResultsResponse
Expand Down Expand Up @@ -323,10 +324,20 @@ def process_nats_pipeline_result(self, job_id: int, result_data: dict, reply_sub
try:
# Save to database (this is the slow operation)
detections_count, classifications_count, captures_count = 0, 0, 0
feature_counts: dict[str, int] = {}
if pipeline_result:
# should never happen since otherwise we could not be processing results here
assert job.pipeline is not None, "Job pipeline is None"
job.pipeline.save_results(results=pipeline_result, job_id=job.pk)
# Only a feature-only save reports boxes that matched no detection and detections
# left without a vector; asking the full save for its created rows would load
# every detection's algorithm.
feature_only = job.pipeline.is_feature_only()
saved = job.pipeline.save_results(results=pipeline_result, job_id=job.pk, return_created=feature_only)
if feature_only and saved:
feature_counts = {
"unmatched": saved.unmatched_detections or 0,
"without_vector": saved.detections_without_vector or 0,
}
job.logger.info(f"Successfully saved results for job {job_id}")

_, t = t(
Expand Down Expand Up @@ -385,6 +396,7 @@ def process_nats_pipeline_result(self, job_id: int, result_data: dict, reply_sub
counts_to_apply = (
(detections_count, classifications_count, captures_count) if is_first_processing else (0, 0, 0)
)
extra_counts = {key: count if is_first_processing else 0 for key, count in feature_counts.items()}
_update_job_progress(
job_id,
"results",
Expand All @@ -393,6 +405,7 @@ def process_nats_pipeline_result(self, job_id: int, result_data: dict, reply_sub
detections=counts_to_apply[0],
classifications=counts_to_apply[1],
captures=counts_to_apply[2],
**extra_counts,
)

# Ack LAST — only after the results-stage SREM and progress commit are
Expand All @@ -408,6 +421,19 @@ def process_nats_pipeline_result(self, job_id: int, result_data: dict, reply_sub
# Celery's autoretry_for handles the transient rather than this broad
# except swallowing it.
raise
except FeatureResultsStoredNothing as e:
# Redelivering would return the same response, and a later job would send the same
# detections again, so record the counts and fail the job instead of retrying.
_update_job_progress(
job_id,
"results",
0,
complete_state=JobState.FAILURE,
unmatched=e.unmatched,
without_vector=e.without_vector,
)
_ack_task_via_nats(reply_subject, job.logger)
_fail_job(job_id, str(e))
except Exception as e:
error = f"Error processing pipeline result for job {job_id}: {e}"
if not acked:
Expand Down Expand Up @@ -502,6 +528,15 @@ async def ack_task():
return False


def _get_stage_param(job, stage: str, key: str) -> int:
"""The integer value of one stage parameter, or 0 when the stage or parameter is absent."""
try:
stage_obj = job.progress.get_stage(stage)
except ValueError:
return 0
return next((param.value or 0 for param in stage_obj.params if param.key == key), 0)


def _get_current_counts_from_job_progress(job, stage: str) -> tuple[int, int, int]:
"""
Get current detections, classifications, and captures counts from job progress.
Expand Down Expand Up @@ -628,6 +663,9 @@ def _update_job_progress(
state_params["detections"] = current_detections + new_detections
state_params["classifications"] = current_classifications + new_classifications
state_params["captures"] = current_captures + new_captures
for key in ("unmatched", "without_vector"):
if key in state_params:
state_params[key] = _get_stage_param(job, stage, key) + state_params[key]

# Don't overwrite a stage with a stale progress value.
# This guards against the race where a slower worker calls _update_job_progress
Expand Down
66 changes: 66 additions & 0 deletions ami/jobs/tests/test_tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -421,6 +421,72 @@ def fail_on_results_stage(self, processed_image_ids, stage, failed_image_ids=Non
self.assertEqual(process_progress.processed, 1)
self.assertEqual(results_progress.processed, 0)

def _feature_only_result(self, image: SourceImage, boxes: list[list[float]]) -> dict:
"""A feature-only result: the given boxes echoed back on one image, each with a vector."""
return PipelineResultsResponse(
pipeline="test-pipeline",
total_time=1.0,
source_images=[SourceImageResponse(id=str(image.pk), url="http://example.com/x.jpg")],
detections=[
{
"source_image_id": str(image.pk),
"bbox": dict(zip(["x1", "y1", "x2", "y2"], box)),
"algorithm": {"name": self.detector.name, "key": self.detector.key},
"timestamp": datetime.datetime.now(),
"embeddings": [
{"algorithm": {"name": self.extractor.name, "key": self.extractor.key}, "features": [0.1] * 8}
],
}
for box in boxes
],
).dict()

def _make_feature_only(self) -> None:
self.detector = Algorithm.objects.create(
name="feature-detector", key="feature-detector", task_type=AlgorithmTaskType.LOCALIZATION
)
self.extractor = Algorithm.objects.create(
name="feature-backbone", key="feature-backbone", task_type=AlgorithmTaskType.EMBEDDING
)
self.pipeline.algorithms.set([self.detector, self.extractor])
for image in self.images:
Detection.objects.create(source_image=image, bbox=[0, 0, 10, 10], detection_algorithm=self.detector)

def _results_param(self, key: str):
self.job.refresh_from_db()
stage = self.job.progress.get_stage("results")
return next((param.value for param in stage.params if param.key == key), None)

@patch("ami.jobs.tasks.TaskQueueManager")
def test_feature_only_results_count_boxes_that_match_no_detection(self, mock_manager_class):
self._setup_mock_nats(mock_manager_class)
self._make_feature_only()
for image in self.images[:2]:
process_nats_pipeline_result(
job_id=self.job.pk,
result_data=self._feature_only_result(image, [[0, 0, 10, 10], [50, 50, 60, 60]]),
reply_subject=f"reply.features.{image.pk}",
)
self.assertEqual(self._results_param("unmatched"), 2)
self.assertEqual(self._results_param("without_vector"), 0)
self.assertNotEqual(self.job.status, JobState.FAILURE.value)

@patch("ami.jobs.tasks._ack_task_via_nats")
@patch("ami.jobs.tasks.TaskQueueManager")
def test_feature_only_batch_matching_no_detection_fails_the_job_and_acks(self, mock_manager_class, mock_ack):
"""Redelivering it would return the same boxes, so the message is acked and the job fails."""
self._setup_mock_nats(mock_manager_class)
self._make_feature_only()
process_nats_pipeline_result(
job_id=self.job.pk,
result_data=self._feature_only_result(self.images[0], [[50, 50, 60, 60], [70, 70, 80, 80]]),
reply_subject="reply.features.none",
)
mock_ack.assert_called_once()
self.assertEqual(self._results_param("unmatched"), 2)
self.assertEqual(self._results_param("without_vector"), 1)
self.assertEqual(self.job.status, JobState.FAILURE.value)

@patch("ami.jobs.tasks.TaskQueueManager")
def test_results_counter_does_not_inflate_on_replay(self, mock_manager_class):
"""
Expand Down
78 changes: 76 additions & 2 deletions ami/main/api/serializers.py
Original file line number Diff line number Diff line change
Expand Up @@ -1182,7 +1182,7 @@ class ClassificationNestedSerializer(ClassificationSerializer):
has_features = serializers.BooleanField(
read_only=True,
allow_null=True,
help_text="Whether a feature embedding was stored for this classification.",
help_text="Whether this classification's algorithm stored a feature vector for its detection.",
)

def get_permissions(self, instance, instance_data):
Expand Down Expand Up @@ -1392,7 +1392,7 @@ class SourceImageSerializer(SourceImageListSerializer):
)
detections_with_features = serializers.IntegerField(
read_only=True,
help_text="Valid detections with at least one classification that stored a feature embedding.",
help_text="Valid detections with a stored feature vector, as an embedding or on a classification.",
)
# file = serializers.ImageField(allow_empty_file=False, use_url=True)

Expand Down Expand Up @@ -1771,6 +1771,10 @@ class OccurrenceSerializer(OccurrenceListSerializer):
event = EventNestedSerializer(read_only=True)
grouping_verified_by = UserNestedSerializer(read_only=True)
grouping_summary = serializers.SerializerMethodField()
grouping_edited_since_verified = serializers.SerializerMethodField(
help_text="Whether the detections changed since the grouping was last confirmed. "
"False while it is confirmed, and when it never was."
)
# first_appearance = TaxonSourceImageNestedSerializer(read_only=True)

class Meta:
Expand All @@ -1788,12 +1792,19 @@ class Meta:
"grouping_verified",
"grouping_verified_at",
"grouping_verified_by",
"grouping_edited_since_verified",
"grouping_summary",
]
read_only_fields = [
"determination_score",
]

def get_grouping_edited_since_verified(self, obj: Occurrence) -> bool:
from ami.main.models_future.history import edited_since_track_complete_review

# Editing the detections withdraws the confirmation, so a confirmed grouping is unchanged.
return obj.grouping_verified_at is None and edited_since_track_complete_review(obj)

@extend_schema_field(OccurrenceFrameSerializer(many=True))
def get_detections(self, obj: Occurrence) -> list[dict]:
frames = frames_from_prefetch(obj)[: self.detections_page_size]
Expand Down Expand Up @@ -2307,6 +2318,52 @@ class OccurrenceGroupingSerializer(serializers.Serializer):
grouping_verified_by = serializers.CharField(allow_null=True)


class HistoryUserSerializer(serializers.Serializer):
"""A person in an occurrence's history: name and picture only, never an email address."""

id = serializers.IntegerField()
name = serializers.CharField()
image = serializers.ImageField(allow_null=True)


class HistoryAlgorithmSerializer(serializers.Serializer):
id = serializers.IntegerField()
name = serializers.CharField()
key = serializers.CharField()


class HistoryJobSerializer(serializers.Serializer):
id = serializers.IntegerField()
name = serializers.CharField()


class HistoryTaxonSerializer(serializers.Serializer):
id = serializers.IntegerField()
name = serializers.CharField()
rank = serializers.CharField()


class OccurrenceHistoryEntrySerializer(serializers.Serializer):
"""One entry of an occurrence's history, newest first. ``type`` says which table it came from."""

type = serializers.ChoiceField(choices=["algorithm_result", "review", "identification", "prediction"])
id = serializers.IntegerField(help_text="Primary key of the row in the table ``type`` names.")
timestamp = serializers.DateTimeField()
subtype = serializers.CharField(
allow_null=True,
help_text="For algorithm results and reviews: tracking, class_masking, size_filter or track_complete.",
)
user = HistoryUserSerializer(allow_null=True)
algorithm = HistoryAlgorithmSerializer(allow_null=True)
job = HistoryJobSerializer(allow_null=True)
taxon = HistoryTaxonSerializer(
allow_null=True, help_text="The identified or predicted taxon, or the determination after a result."
)
taxon_before = HistoryTaxonSerializer(allow_null=True, help_text="The determination before a result.")
score = serializers.FloatField(allow_null=True)
payload = serializers.JSONField(help_text="Details that depend on the type and subtype.")


class OccurrencePathCaptureSerializer(serializers.Serializer):
"""The capture one frame of a path was measured against."""

Expand Down Expand Up @@ -2504,3 +2561,20 @@ class CaptureMatchesResponseSerializer(serializers.Serializer):
"and then every score is null.",
)
detections = CaptureMatchSerializer(many=True)


class SessionFeatureExtractorSerializer(serializers.Serializer):
"""A feature extractor with vectors stored for a session's detections."""

id = serializers.IntegerField(source="algorithm.pk")
name = serializers.CharField(source="algorithm.name")
key = serializers.CharField(source="algorithm.key")
task_type = serializers.CharField(source="algorithm.task_type", allow_null=True)
embedding_dimensions = serializers.IntegerField(source="algorithm.embedding_dimensions", allow_null=True)
embeddings_count = serializers.IntegerField(help_text="Detections with a stored embedding from it.")
classification_vectors_count = serializers.IntegerField(
help_text="Classifications from it that carry a vector (data processed before embeddings were stored)."
)
is_default = serializers.BooleanField(
help_text="Whether tracking compares this extractor's vectors when none is chosen."
)
Loading