From 74869bbfe2a8459cfa4047281f112d4b56a03a85 Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Wed, 23 Sep 2026 06:54:22 -0700 Subject: [PATCH 1/2] fix(jobs): stop creating a job from asking the result backend for its status Enqueuing a job read the status of the task it had just sent from the Celery result backend, inside the request that creates the job. When that backend connection had gone idle and been reset, the read raised ConnectionResetError and the request failed with a 500, which was seen repeatedly when starting tracking and export jobs. A task that was just sent is always PENDING, so the job is now marked PENDING directly and creating a job no longer talks to the result backend. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01C7Xf6VPbwWtTumhjjF15g8 --- ami/jobs/models.py | 4 +++- ami/jobs/tests/test_jobs.py | 17 +++++++++++++++++ 2 files changed, 20 insertions(+), 1 deletion(-) diff --git a/ami/jobs/models.py b/ami/jobs/models.py index ff65f31f2..ed772c9fa 100644 --- a/ami/jobs/models.py +++ b/ami/jobs/models.py @@ -1135,7 +1135,9 @@ def send_task(): self.started_at = None self.finished_at = None self.scheduled_at = datetime.datetime.now() - self.status = AsyncResult(task_id).status + # A task that was just sent is PENDING by definition. Asking the result backend + # adds nothing and fails the request when its idle connection has been reset. + self.status = JobState.PENDING self.update_progress(save=False) self.save() diff --git a/ami/jobs/tests/test_jobs.py b/ami/jobs/tests/test_jobs.py index 00b7934a7..31dd8b5e4 100644 --- a/ami/jobs/tests/test_jobs.py +++ b/ami/jobs/tests/test_jobs.py @@ -1,6 +1,7 @@ # from rich import print import logging from typing import Any +from unittest.mock import patch from django.test import TestCase from guardian.shortcuts import assign_perm @@ -1608,6 +1609,22 @@ def test_regroup_job_failure_propagates(self): job.run() +class TestJobEnqueue(TestCase): + def test_enqueue_marks_the_job_pending_without_asking_the_result_backend(self): + """Creating a job must not depend on the result backend, whose idle connection + can be reset and turn the request into a 500.""" + project = Project.objects.create(name="Enqueue Project") + job = Job.objects.create(name="Enqueue test", project=project, job_type_key=RegroupEventsJob.key) + + with patch("ami.jobs.models.AsyncResult") as async_result: + job.enqueue() + + async_result.assert_not_called() + job.refresh_from_db() + self.assertEqual(job.status, JobState.PENDING.value) + self.assertIsNotNone(job.task_id) + + class TestDataStorageSyncJobIncludesRegroupStage(TestCase): """ The sync Job now runs grouping as an explicit second stage so logs land on From 24134f6d9d3efe7402010a0d800f88b5da03b508 Mon Sep 17 00:00:00 2001 From: Michael Bunsen Date: Mon, 28 Sep 2026 21:00:40 -0700 Subject: [PATCH 2/2] fix(jobs): save the pending state before sending the job's task Outside a transaction on_commit runs its callback at once, so a fast worker could save STARTED before enqueue() wrote PENDING over it. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01C7Xf6VPbwWtTumhjjF15g8 --- ami/jobs/models.py | 6 ++++-- ami/jobs/tests/test_jobs.py | 18 ++++++++++++++++++ 2 files changed, 22 insertions(+), 2 deletions(-) diff --git a/ami/jobs/models.py b/ami/jobs/models.py index ed772c9fa..590738841 100644 --- a/ami/jobs/models.py +++ b/ami/jobs/models.py @@ -1130,16 +1130,18 @@ def enqueue(self): def send_task(): run_job.apply_async(kwargs={"job_id": self.pk}, task_id=task_id) - transaction.on_commit(send_task) self.task_id = task_id self.started_at = None self.finished_at = None self.scheduled_at = datetime.datetime.now() - # A task that was just sent is PENDING by definition. Asking the result backend + # A task that is about to be sent is PENDING by definition. Asking the result backend # adds nothing and fails the request when its idle connection has been reset. self.status = JobState.PENDING self.update_progress(save=False) self.save() + # Dispatch only after PENDING is saved: outside a transaction on_commit runs at once, + # and a fast worker's STARTED would otherwise be overwritten by the save above. + transaction.on_commit(send_task) def setup(self, save=True): """ diff --git a/ami/jobs/tests/test_jobs.py b/ami/jobs/tests/test_jobs.py index 31dd8b5e4..dc0326a3e 100644 --- a/ami/jobs/tests/test_jobs.py +++ b/ami/jobs/tests/test_jobs.py @@ -1624,6 +1624,24 @@ def test_enqueue_marks_the_job_pending_without_asking_the_result_backend(self): self.assertEqual(job.status, JobState.PENDING.value) self.assertIsNotNone(job.task_id) + def test_enqueue_saves_pending_before_the_task_is_sent(self): + """The task must be sent only after PENDING is saved, or a fast worker's STARTED + state could be overwritten by the enqueue save.""" + project = Project.objects.create(name="Enqueue Project") + job = Job.objects.create(name="Enqueue test", project=project, job_type_key=RegroupEventsJob.key) + status_at_dispatch = [] + + def record_status(*args, **kwargs): + status_at_dispatch.append(Job.objects.get(pk=job.pk).status) + + # Outside a transaction on_commit runs the callback at once; simulate that here. + with patch("ami.jobs.models.run_job.apply_async", side_effect=record_status), patch( + "ami.jobs.models.transaction.on_commit", side_effect=lambda fn: fn() + ): + job.enqueue() + + self.assertEqual(status_at_dispatch, [JobState.PENDING.value]) + class TestDataStorageSyncJobIncludesRegroupStage(TestCase): """