Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
f2ee950
feat(tracking): link detections into chains, store occurrence statist…
mihow Oct 4, 2026
48e0205
feat(exports): add a CSV export with one row per detection of each oc…
mihow Oct 4, 2026
b5c5add
feat(post-processing): generate admin action forms from a task's conf…
mihow Oct 4, 2026
8711c7b
feat(tracking): run occurrence tracking as a post-processing task fro…
mihow Oct 4, 2026
7b389d3
docs(tracking): say that re-tracking a session adds links but never u…
mihow Oct 4, 2026
c3c9b6a
refactor(exports): drop the dedicated tracks export
mihow Oct 5, 2026
334f7f4
refactor(tracking): drop the stored occurrence statistics
mihow Oct 5, 2026
249c3c0
refactor(tracking): give tracking one home under ami/ml with pure cod…
mihow Oct 5, 2026
d503ed3
chore(migrations): renumber the next_detection migration after the al…
mihow Oct 6, 2026
4d55301
feat(tracking): record a tracking result for each occurrence a run links
mihow Oct 6, 2026
a3f5a51
feat(ui): show the tracking result on an occurrence's history
mihow Oct 6, 2026
1245868
fix(migrations): add the detection next-detection link without blocki…
mihow Oct 6, 2026
c486d83
fix(regroup): look for occurrences split by a regroup only among the …
mihow Oct 6, 2026
940749e
refactor(tracking): drop the tracking classification, record link cos…
mihow Oct 6, 2026
c1cb15d
feat(ui): show the reworked tracking figures on an occurrence's history
mihow Oct 6, 2026
78437c1
fix(tracking): save job progress while a session is being matched, an…
mihow Oct 6, 2026
5283b0f
feat(tracking): label the session list setting so the history shows i…
mihow Oct 6, 2026
2eb150a
test(ui): give the tracking history fixture a value and no score
mihow Oct 6, 2026
26a0d83
fix(tracking): never delete an occurrence that still holds detections…
mihow Oct 6, 2026
2353d77
refactor(tracking): follow the result framework's kind and settings d…
mihow Oct 6, 2026
0072d30
refactor(tracking): declare the tracking result model and drop the cu…
mihow Oct 6, 2026
e989d46
test(tracking): build tracking test fixtures once per class
mihow Oct 6, 2026
caadb93
refactor(tracking): let the result framework derive the tracking valu…
mihow Oct 7, 2026
d68a2ad
perf(main): stop update_occurrence_determination from running the loo…
mihow Oct 7, 2026
ed114d9
feat(tracking): write a session in bulk, record the grouping for undo…
mihow Oct 7, 2026
049ce4b
fix(tracking): skip the regroup split where nothing was tracked, and …
mihow Oct 7, 2026
d4bb07b
fix(migrations): give up after 10 s rather than queue behind a long q…
mihow Oct 7, 2026
8b387eb
refactor(tracking): reuse the admin's error mapping and drop the card…
mihow Oct 7, 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
22 changes: 21 additions & 1 deletion ami/main/admin.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,11 @@
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.small_size_filter_form import SmallSizeFilterActionForm
from ami.ml.post_processing.admin.tracking_actions import build_tracking_jobs_for_events
from ami.ml.post_processing.admin.tracking_form import TrackingActionForm
from ami.ml.post_processing.class_masking import ClassMaskingTask
from ami.ml.post_processing.small_size_filter import SmallSizeFilterTask
from ami.ml.post_processing.tracking import TrackingTask
from ami.ml.tasks import remove_duplicate_classifications

from .models import (
Expand Down Expand Up @@ -314,8 +317,16 @@ def update_calculated_fields(self, request: HttpRequest, queryset: QuerySet[Even
update_calculated_fields_for_events(qs=queryset)
self.message_user(request, f"Updated {queryset.count()} events.")

# One Job per project, since a Job belongs to a single project and the changelist can span several.
run_tracking = make_post_processing_action(
TrackingTask,
TrackingActionForm,
build_jobs=build_tracking_jobs_for_events,
description="Run Occurrence tracking on the selected sessions (async)",
)

list_filter = ("deployment", "project", "start")
actions = [update_calculated_fields]
actions = [update_calculated_fields, run_tracking]


@admin.register(SourceImage)
Expand Down Expand Up @@ -861,11 +872,20 @@ def populate_collection_async(self, request: HttpRequest, queryset: QuerySet[Sou
f"Post-processing: {task_cls.name} on Capture Set {collection.pk}"
),
)
run_tracking = make_post_processing_action(
TrackingTask,
TrackingActionForm,
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_tracking,
]

# Hide images many-to-many field from form. This would list all source images in the database.
Expand Down
45 changes: 45 additions & 0 deletions ami/main/migrations/0100_detection_next_detection.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,45 @@
import django.db.models.deletion
from django.db import migrations, models


class Migration(migrations.Migration):
"""Add the column that links a detection to the one that follows it in a tracking sequence.

Django would add this column together with a unique constraint and a foreign key in one
statement, which takes a lock on the large detection table that blocks reads while the
unique index is built. Here the column is added on its own: a nullable column without a
default is a catalogue change that does not scan the table. Migration 0101 adds the unique
constraint and the foreign key afterwards without blocking readers or writers.
"""

dependencies = [
("main", "0099_classification_algorithm_result_index"),
]

operations = [
# Adding the column takes a brief exclusive lock. Give up rather than queue behind a long query on
# the table, which would block every other query until it ends; rerun the migration if it gives up.
migrations.RunSQL(sql="SET LOCAL lock_timeout = '10s';", reverse_sql=migrations.RunSQL.noop),
migrations.SeparateDatabaseAndState(
state_operations=[
migrations.AddField(
model_name="detection",
name="next_detection",
field=models.OneToOneField(
blank=True,
help_text="The detection that follows this one in the tracking sequence.",
null=True,
on_delete=django.db.models.deletion.SET_NULL,
related_name="previous_detection",
to="main.detection",
),
),
],
database_operations=[
migrations.RunSQL(
sql='ALTER TABLE "main_detection" ADD COLUMN "next_detection_id" bigint NULL;',
reverse_sql='ALTER TABLE "main_detection" DROP COLUMN "next_detection_id";',
),
],
),
]
72 changes: 72 additions & 0 deletions ami/main/migrations/0101_detection_next_detection_constraints.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
from django.db import migrations


class Migration(migrations.Migration):
"""Give ``Detection.next_detection`` the unique constraint and foreign key that Django would have created.

Both are built so that neither blocks reads or writes on the large detection table. The
unique index is built CONCURRENTLY, which needs a non-atomic migration, and is then attached
as a constraint, which is a catalogue change. The foreign key is added NOT VALID, so existing
rows are not checked while a strong lock is held, and validated afterwards, which takes only a
light lock. See 0093 for why the statement timeout is cleared and restored around the build.

The constraint names match the ones Django generates, so later AlterField migrations find them.
If the index build is interrupted it leaves an invalid index of the same name, and if a lock times out the
index already exists; drop it before retrying.
"""

atomic = False

dependencies = [
("main", "0100_detection_next_detection"),
]

operations = [
migrations.RunSQL(
sql="SET statement_timeout = 0;",
reverse_sql=migrations.RunSQL.noop,
),
migrations.RunSQL(
sql=(
'CREATE UNIQUE INDEX CONCURRENTLY "main_detection_next_detection_id_key" '
'ON "main_detection" ("next_detection_id");'
),
reverse_sql=migrations.RunSQL.noop,
),
# The two ALTER TABLE statements below take brief strong locks. Give up rather than queue behind a long
# query on the table, which would block every other query until it ends.
migrations.RunSQL(sql="SET lock_timeout = '10s';", reverse_sql=migrations.RunSQL.noop),
migrations.RunSQL(
sql=(
'ALTER TABLE "main_detection" ADD CONSTRAINT "main_detection_next_detection_id_key" '
'UNIQUE USING INDEX "main_detection_next_detection_id_key";'
),
reverse_sql='ALTER TABLE "main_detection" DROP CONSTRAINT "main_detection_next_detection_id_key";',
),
migrations.RunSQL(
sql=(
'ALTER TABLE "main_detection" ADD CONSTRAINT "main_detection_next_detection_id_f0201e13_fk_main_detection_id" '
'FOREIGN KEY ("next_detection_id") REFERENCES "main_detection" ("id") '
"DEFERRABLE INITIALLY DEFERRED NOT VALID;"
),
reverse_sql=(
'ALTER TABLE "main_detection" '
'DROP CONSTRAINT "main_detection_next_detection_id_f0201e13_fk_main_detection_id";'
),
),
migrations.RunSQL(
sql=(
'ALTER TABLE "main_detection" VALIDATE CONSTRAINT '
'"main_detection_next_detection_id_f0201e13_fk_main_detection_id";'
),
reverse_sql=migrations.RunSQL.noop,
),
migrations.RunSQL(
sql="RESET lock_timeout;",
reverse_sql=migrations.RunSQL.noop,
),
migrations.RunSQL(
sql="RESET statement_timeout;",
reverse_sql=migrations.RunSQL.noop,
),
]
92 changes: 85 additions & 7 deletions ami/main/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -1441,6 +1441,20 @@ def update_calculated_fields_for_events(
return to_update


def update_calculated_fields_for_sessions_and_stations(event_ids: typing.Iterable[int | None]) -> None:
"""Refresh the cached counts of these sessions and of the stations they belong to.

Call once after occurrences are created, merged or split, which neither the occurrence nor the
detection saves do. The station refresh scans each whole station, so call it from a background job.
"""
pks = sorted({pk for pk in event_ids if pk is not None})
if not pks:
return
update_calculated_fields_for_events(pks=pks)
for deployment in Deployment.objects.filter(events__pk__in=pks).distinct():
deployment.update_calculated_fields(save=True)


def audit_event_lengths(deployment: Deployment):
logger.info("Checking for unusual event durations")

Expand Down Expand Up @@ -1648,6 +1662,8 @@ def _group_images_into_events_locked(
f"Done grouping {len(image_timestamps)} captures into {len(events)} events " f"for deployment {deployment}"
)

occurrences_split_count = _split_occurrences_at_session_boundaries(job, touched_event_pks)

# Realign Occurrence.event_id with each occurrence's detections' current
# source_image.event_id. Occurrences are bound to an event once at creation
# time (Detection.associate_new_occurrence and Pipeline.save_results both
Expand Down Expand Up @@ -1728,6 +1744,7 @@ def _group_images_into_events_locked(
"Events created": events_created_count,
"Events touched": len(touched_event_pks),
"Empty events deleted": events_deleted_empty,
"Occurrences split at a session boundary": occurrences_split_count,
"Duplicate timestamps": duplicate_timestamp_count,
"Ungrouped captures": ungrouped_captures_count,
"Captures missing timestamp": no_timestamp_captures_count,
Expand All @@ -1740,6 +1757,62 @@ def _group_images_into_events_locked(
return events


def _split_occurrences_at_session_boundaries(job: "Job | None", event_pks: set[int]) -> int:
"""Split every occurrence whose detections now span several sessions, among the sessions a regroup touched.

An occurrence is expected to belong to one session, so a regroup that draws a session
boundary through it leaves one piece per session. Tracking is what merges detections of several
captures into one occurrence, so when no detection of these sessions has a tracking link the search is
skipped after one indexed query; an occurrence grouped some other way, with no links, is then not split.
The search itself reads every occurrence of the touched sessions. Returns how many occurrences were split.
"""
from ami.ml.post_processing.tracking.sessions import lock_sessions, split_at_session_boundaries

if not Detection.objects.filter(source_image__event_id__in=event_pks, next_detection__isnull=False).exists():
return 0

def find_spanning_ids() -> list[int]:
capture_ids = list(SourceImage.objects.filter(event_id__in=event_pks).values_list("pk", flat=True))
touched_occurrence_ids = list(
Detection.objects.filter(source_image_id__in=capture_ids, occurrence__isnull=False)
.values_list("occurrence_id", flat=True)
.distinct()
)
return list(
Detection.objects.valid()
.filter(occurrence_id__in=touched_occurrence_ids)
.values("occurrence_id")
.annotate(sessions=models.Count("source_image__event", distinct=True))
.filter(sessions__gt=1)
.values_list("occurrence_id", flat=True)
)

candidate_ids = find_spanning_ids()
if not candidate_ids:
return 0
split_count = 0
# Holding the sessions' locks while the occurrences are found and split makes a tracking run
# on one of these sessions finish first, or wait for the split.
with transaction.atomic():
lock_sessions(
list(
SourceImage.objects.filter(detections__occurrence_id__in=candidate_ids)
.values_list("event_id", flat=True)
.distinct()
)
)
for occurrence in Occurrence.objects.filter(pk__in=find_spanning_ids()).order_by("pk"):
pieces = split_at_session_boundaries(occurrence)
if not pieces:
continue
split_count += 1
(job.logger if job else logger).info(
f"Split occurrence {occurrence.pk} at a session boundary; "
f"new occurrence(s) {[piece.pk for piece in pieces]} hold the later sessions."
)
return split_count


def deployment_events_need_update(deployment: Deployment) -> bool:
"""
Returns True if there are any SourceImages in the deployment
Expand Down Expand Up @@ -3239,6 +3312,15 @@ class Detection(BaseModel):

similarity_vector = models.JSONField(null=True, blank=True)

next_detection = models.OneToOneField(
"self",
on_delete=models.SET_NULL,
null=True,
blank=True,
related_name="previous_detection",
help_text="The detection that follows this one in the tracking sequence.",
)

# For type hints
classifications: models.QuerySet["Classification"]
source_image_id: int
Expand Down Expand Up @@ -3896,13 +3978,9 @@ def update_occurrence_determination(
"""
needs_update = False

# Invalidate the cached properties so they will be re-calculated
if hasattr(occurrence, "best_identification"):
del occurrence.best_identification
if hasattr(occurrence, "best_prediction"):
del occurrence.best_prediction
if hasattr(occurrence, "best_identification"):
del occurrence.best_identification
# Clear the cached properties so they are recalculated. ``hasattr`` would run their queries first.
occurrence.__dict__.pop("best_identification", None)
occurrence.__dict__.pop("best_prediction", None)

current_determination = (
current_determination
Expand Down
Loading