Skip to content

Commit 8a717f4

Browse files
committed
Fix duplicate recurring enqueues across scheduler instances
1 parent 019557b commit 8a717f4

4 files changed

Lines changed: 57 additions & 5 deletions

File tree

CHANGELOG.md

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,13 @@
11
# Changelog
22

3+
## Unreleased
4+
5+
**Fixed:**
6+
7+
- Avoid duplicate recurring-task enqueues when multiple schedulers race on the
8+
same `run_at`. We now create recurring execution records atomically and skip
9+
already-recorded runs, matching Solid Queue's behavior.
10+
311
## v0.2.0 - 2026-04-18
412

513
**Fixed:**

steady_queue/models/recurring_execution.py

Lines changed: 8 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
from django.db import models
1+
from django.db import IntegrityError, models
22
from django.tasks import TaskResult
33

44
from .execution import Execution, ExecutionQuerySet
@@ -9,9 +9,10 @@ def clearable(self):
99
return self.filter(job__isnull=True)
1010

1111
def record(self, task_result: TaskResult, task, run_at):
12-
self.update_or_create(
13-
task=task, run_at=run_at, defaults={"job_id": task_result.id}
14-
)
12+
try:
13+
self.create(task=task, run_at=run_at, job_id=task_result.id)
14+
except IntegrityError as e:
15+
raise self.model.AlreadyRecorded from e
1516

1617
def clear_in_batches(self, batch_size=500):
1718
while True:
@@ -21,6 +22,9 @@ def clear_in_batches(self, batch_size=500):
2122

2223

2324
class RecurringExecution(Execution):
25+
class AlreadyRecorded(Exception):
26+
pass
27+
2428
class Meta:
2529
verbose_name = "recurring execution"
2630
verbose_name_plural = "recurring executions"

steady_queue/models/recurring_task.py

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -97,7 +97,10 @@ def last_enqueued_time(self) -> datetime:
9797
return self.recurring_executions.max("run_at")
9898

9999
def enqueue(self, run_at: datetime):
100-
return self.enqueue_and_record(run_at)
100+
try:
101+
return self.enqueue_and_record(run_at)
102+
except RecurringExecution.AlreadyRecorded:
103+
return False
101104

102105
def enqueue_and_record(self, run_at: datetime):
103106
with transaction.atomic(using=self._state.db):
@@ -108,6 +111,7 @@ def enqueue_and_record(self, run_at: datetime):
108111
queue_name=self.queue_name, priority=self.priority
109112
).enqueue(*args, **kwargs)
110113
RecurringExecution.objects.record(task_result, self, run_at)
114+
return task_result
111115

112116
@property
113117
def previous_time(self) -> datetime:

tests/test_recurring_execution.py

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,36 @@
1+
from django.test import TestCase
2+
from django.utils import timezone
3+
4+
from steady_queue.configuration import Configuration
5+
from steady_queue.models import Job, RecurringExecution, RecurringTask
6+
from tests.dummy.tasks import dummy_task
7+
8+
9+
class RecurringExecutionTestCase(TestCase):
10+
def test_enqueue_same_run_at_only_enqueues_once(self):
11+
recurring_task = RecurringTask.from_configuration(
12+
Configuration.RecurringTask(
13+
key="test-recurring-task",
14+
class_name="tests.dummy.tasks.dummy_task",
15+
schedule="* * * * *",
16+
arguments=dummy_task.serialize([], {}),
17+
queue_name="default",
18+
priority=0,
19+
)
20+
)
21+
recurring_task.save()
22+
23+
run_at = timezone.now().replace(second=0, microsecond=0)
24+
25+
first_result = recurring_task.enqueue(run_at=run_at)
26+
second_result = recurring_task.enqueue(run_at=run_at)
27+
28+
self.assertNotEqual(first_result, False)
29+
self.assertFalse(second_result)
30+
31+
self.assertEqual(Job.objects.count(), 1)
32+
self.assertEqual(RecurringExecution.objects.count(), 1)
33+
34+
recurring_execution = RecurringExecution.objects.get()
35+
self.assertEqual(recurring_execution.task_id, recurring_task.key)
36+
self.assertEqual(recurring_execution.run_at, run_at)

0 commit comments

Comments
 (0)