Skip to content
Open
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
8 changes: 6 additions & 2 deletions ami/jobs/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -1130,14 +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()
self.status = AsyncResult(task_id).status
# 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
Comment thread
mihow marked this conversation as resolved.
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):
"""
Expand Down
35 changes: 35 additions & 0 deletions ami/jobs/tests/test_jobs.py
Original file line number Diff line number Diff line change
@@ -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
Expand Down Expand Up @@ -1608,6 +1609,40 @@ 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)

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):
"""
The sync Job now runs grouping as an explicit second stage so logs land on
Expand Down
Loading