Skip to content

Commit 441a610

Browse files
committed
fix: schedule tasks with on_commit
1 parent 82d6161 commit 441a610

4 files changed

Lines changed: 25 additions & 19 deletions

File tree

share/models/index_backfill.py

Lines changed: 5 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -98,14 +98,11 @@ def pls_start(self, index_strategy):
9898
locked_self.strategy_checksum = _current_checksum
9999
locked_self.backfill_status = IndexBackfill.INITIAL
100100
locked_self.__update_error(None)
101-
try:
102-
task__schedule_index_backfill.apply_async((locked_self.pk,))
103-
except Exception as error:
104-
locked_self.__update_error(error)
105-
else:
106-
locked_self.backfill_status = IndexBackfill.WAITING
107-
finally:
108-
locked_self.save()
101+
locked_self.backfill_status = IndexBackfill.WAITING
102+
locked_self.save()
103+
transaction.on_commit(
104+
lambda: task__schedule_index_backfill.apply_async((locked_self.pk,))
105+
)
109106

110107
def pls_note_scheduling_has_begun(self):
111108
with self.mutex() as locked_self:

tests/share/search/test_index_backfill.py

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -20,12 +20,15 @@ def index_backfill(self, fake_strategy):
2020
index_strategy_name=fake_strategy.strategy_name,
2121
)
2222

23-
def test_happypath(self, index_backfill: IndexBackfill, fake_strategy):
23+
def test_happypath(self, index_backfill: IndexBackfill, fake_strategy, django_capture_on_commit_callbacks):
2424
assert index_backfill.backfill_status == IndexBackfill.INITIAL
2525
assert index_backfill.strategy_checksum == ''
26-
with mock.patch('share.models.index_backfill.task__schedule_index_backfill') as mock_task:
26+
with (
27+
mock.patch('share.models.index_backfill.task__schedule_index_backfill') as mock_task,
28+
django_capture_on_commit_callbacks(execute=True),
29+
):
2730
index_backfill.pls_start(fake_strategy)
28-
mock_task.apply_async.assert_called_once_with((index_backfill.pk,))
31+
mock_task.apply_async.assert_called_once_with((index_backfill.pk,))
2932
assert index_backfill.backfill_status == IndexBackfill.WAITING
3033
assert index_backfill.strategy_checksum == 'foo_bar'
3134
index_backfill.pls_note_scheduling_has_begun()

tests/trove/digestive_tract/test_expel.py

Lines changed: 12 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -61,11 +61,13 @@ def test_setup(self):
6161
def test_expel(self):
6262
with mock.patch('trove.digestive_tract.expel_suid') as _mock_expel_suid:
6363
_user = self.suid_1.source_config.source.user
64-
digestive_tract.expel(from_user=_user, record_identifier=self.suid_1.identifier)
64+
with self.captureOnCommitCallbacks(execute=True):
65+
digestive_tract.expel(from_user=_user, record_identifier=self.suid_1.identifier)
6566
_mock_expel_suid.assert_called_once_with(self.suid_1)
6667

6768
def test_expel_suid(self):
68-
digestive_tract.expel_suid(self.suid_1)
69+
with self.captureOnCommitCallbacks(execute=True):
70+
digestive_tract.expel_suid(self.suid_1)
6971
self.indexcard_1.refresh_from_db()
7072
self.indexcard_2.refresh_from_db()
7173
self.assertIsNotNone(self.indexcard_1.deleted)
@@ -85,7 +87,8 @@ def test_expel_suid(self):
8587
self.mock_derive_task.delay.assert_not_called()
8688

8789
def test_expel_supplementary_suid(self):
88-
digestive_tract.expel_suid(self.supp_suid)
90+
with self.captureOnCommitCallbacks(execute=True):
91+
digestive_tract.expel_suid(self.supp_suid)
8992
self.indexcard_1.refresh_from_db()
9093
self.indexcard_2.refresh_from_db()
9194
self.assertIsNone(self.indexcard_1.deleted)
@@ -105,15 +108,17 @@ def test_expel_supplementary_suid(self):
105108

106109
def test_expel_expired_task(self):
107110
with mock.patch('trove.digestive_tract.expel_expired_data') as _mock_expel_expired:
108-
digestive_tract.task__expel_expired_data.apply()
111+
with self.captureOnCommitCallbacks(execute=True):
112+
digestive_tract.task__expel_expired_data.apply()
109113
_mock_expel_expired.assert_called_once_with(datetime.date.today())
110114

111115
def test_expel_expired(self):
112116
_today = datetime.date.today()
113117
_latest = self.indexcard_2.latest_resource_description
114118
_latest.expiration_date = _today
115119
_latest.save()
116-
digestive_tract.expel_expired_data(_today)
120+
with self.captureOnCommitCallbacks(execute=True):
121+
digestive_tract.expel_expired_data(_today)
117122
self.indexcard_1.refresh_from_db()
118123
self.indexcard_2.refresh_from_db()
119124
self.assertIsNone(self.indexcard_1.deleted)
@@ -136,7 +141,8 @@ def test_expel_expired_supplement(self):
136141
_today = datetime.date.today()
137142
self.supp.expiration_date = _today
138143
self.supp.save()
139-
digestive_tract.expel_expired_data(_today)
144+
with self.captureOnCommitCallbacks(execute=True):
145+
digestive_tract.expel_expired_data(_today)
140146
self.indexcard_1.refresh_from_db()
141147
self.indexcard_2.refresh_from_db()
142148
self.assertIsNone(self.indexcard_1.deleted)

trove/digestive_tract.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,7 @@ def ingest(
6565
expiration_date=expiration_date,
6666
)
6767
for _card in _extracted_cards:
68-
task__derive.delay(_card.pk, urgent=urgent)
68+
transaction.on_commit(lambda: task__derive.delay(_card.pk, urgent=urgent))
6969

7070

7171
@transaction.atomic
@@ -249,7 +249,7 @@ def _expel_supplementary_descriptions(supplementary_rdf_queryset: QuerySet[trove
249249
_affected_indexcards.add(_supplement.indexcard)
250250
_supplement.delete()
251251
for _indexcard in _affected_indexcards:
252-
task__derive.delay(_indexcard.pk)
252+
transaction.on_commit(lambda: task__derive.delay(_indexcard.pk))
253253

254254

255255
### BEGIN celery tasks

0 commit comments

Comments
 (0)