Skip to content

Commit e67fbd0

Browse files
committed
refactor self-ref order fallback to use model-scoped db updates
1 parent 40be75d commit e67fbd0

2 files changed

Lines changed: 78 additions & 39 deletions

File tree

morango/sync/operations.py

Lines changed: 30 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,10 @@
88
from django.core import exceptions
99
from django.db import connection
1010
from django.db.models import CharField
11+
from django.db.models import Exists
12+
from django.db.models import OuterRef
1113
from django.db.models import Q
14+
from django.db.models import Subquery
1215
from django.db.models import signals
1316
from django.db.utils import OperationalError
1417
from django.utils import timezone
@@ -680,49 +683,38 @@ def _queue_into_buffer_v2(transfersession, chunk_size=200):
680683
)
681684

682685

686+
def _update_legacy_self_ref_order_for_model(queryset):
687+
# root nodes set the _self_ref_order to 0
688+
queryset.filter(_self_ref_fk="").exclude(_self_ref_order=0).update(_self_ref_order=0)
689+
# reset the _self_ref_order to None for all records that have a parent
690+
queryset.exclude(_self_ref_fk="").exclude(_self_ref_order=None).update(
691+
_self_ref_order=None
692+
)
693+
694+
parent = Store.objects.filter(
695+
id=OuterRef("_self_ref_fk"),
696+
_self_ref_order__isnull=False,
697+
)
698+
parent_order = parent.values("_self_ref_order")[:1]
699+
pending = queryset.exclude(_self_ref_fk="").filter(_self_ref_order=None)
700+
701+
while pending.filter(Exists(parent)).update(_self_ref_order=Subquery(parent_order)):
702+
pass
703+
704+
683705
def _update_legacy_self_ref_order(transfer_session):
706+
profile = transfer_session.sync_session.profile
684707
transferred_store_records = Store.objects.filter(
685-
last_transfer_session_id=transfer_session.id
708+
last_transfer_session_id=transfer_session.id,
709+
profile=profile,
686710
)
687-
records_by_id = {record.id: record for record in transferred_store_records}
688-
records_by_model = defaultdict(list)
689-
690-
for record in records_by_id.values():
691-
try:
692-
Model = syncable_models.get_model(record.profile, record.model_name)
693-
except KeyError:
694-
continue
695711

712+
for Model in syncable_models.get_models(profile):
713+
queryset = transferred_store_records.filter(model_name=Model.morango_model_name)
696714
if self_referential_fk(Model):
697-
records_by_model[(record.profile, record.model_name)].append(record)
698-
elif record._self_ref_order is not None:
699-
record._self_ref_order = None
700-
record.save(update_fields=["_self_ref_order"])
701-
702-
for records in records_by_model.values():
703-
pending_record_ids = set(record.id for record in records)
704-
while pending_record_ids:
705-
updated = False
706-
for record_id in list(pending_record_ids):
707-
record = records_by_id[record_id]
708-
if not record._self_ref_fk:
709-
order = 0
710-
else:
711-
parent = records_by_id.get(record._self_ref_fk)
712-
if parent is not None and parent.id in pending_record_ids:
713-
continue
714-
order = Store.objects.filter(id=record._self_ref_fk).values_list(
715-
"_self_ref_order", flat=True
716-
).first()
717-
718-
if record._self_ref_order != order:
719-
record._self_ref_order = order
720-
record.save(update_fields=["_self_ref_order"])
721-
pending_record_ids.remove(record_id)
722-
updated = True
723-
724-
if not updated:
725-
break
715+
_update_legacy_self_ref_order_for_model(queryset)
716+
else:
717+
queryset.exclude(_self_ref_order=None).update(_self_ref_order=None)
726718

727719

728720
def _dequeue_into_store(transfer_session, fsic, v2_format=False, self_ref_order=True):

tests/testapp/tests/sync/test_operations.py

Lines changed: 48 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@
3737
_deserialize_from_store,
3838
_queue_into_buffer_v1,
3939
_queue_into_buffer_v2,
40+
_update_legacy_self_ref_order,
4041
)
4142
from morango.sync.syncsession import TransferClient
4243

@@ -654,7 +655,7 @@ def setUp(self):
654655
conn.server_info = dict(capabilities=[])
655656
self.data["mc"] = MorangoProfileController("facilitydata")
656657
session = SyncSession.objects.create(
657-
id=uuid.uuid4().hex, profile="", last_activity_timestamp=timezone.now()
658+
id=uuid.uuid4().hex, profile="facilitydata", last_activity_timestamp=timezone.now()
658659
)
659660
self.transfer_session = TransferSession.objects.create(
660661
id=uuid.uuid4().hex,
@@ -684,6 +685,21 @@ def assert_store_records_not_tagged_with_last_session(self, store_ids):
684685
except Store.DoesNotExist:
685686
pass
686687

688+
def _make_transferred_store(self, **kwargs):
689+
defaults = {
690+
"id": uuid.uuid4().hex,
691+
"serialized": "{}",
692+
"last_saved_instance": self.current_id.id,
693+
"last_saved_counter": 1,
694+
"model_name": "facility",
695+
"profile": "facilitydata",
696+
"partition": uuid.uuid4().hex,
697+
"source_id": uuid.uuid4().hex,
698+
"last_transfer_session_id": self.transfer_session.id,
699+
}
700+
defaults.update(kwargs)
701+
return Store.objects.create(**defaults)
702+
687703
def test_dequeuing_sets_last_session(self):
688704
store_ids = [self.data[key] for key in ["model2", "model3", "model4", "model5", "model7"]]
689705
self.assert_store_records_not_tagged_with_last_session(store_ids)
@@ -1009,6 +1025,37 @@ def test_dequeue_into_store__self_ref_order_fallback_for_missing_capability(self
10091025
self.assertEqual(Store.objects.get(id=self.data["model3"])._self_ref_order, 0)
10101026
self.assertEqual(Store.objects.get(id=self.data["model4"])._self_ref_order, 0)
10111027

1028+
def test_update_legacy_self_ref_order_nulls_non_self_ref_models(self):
1029+
store = self._make_transferred_store(
1030+
model_name=SummaryLog.morango_model_name,
1031+
_self_ref_order=3,
1032+
)
1033+
1034+
_update_legacy_self_ref_order(self.transfer_session)
1035+
1036+
store.refresh_from_db()
1037+
self.assertIsNone(store._self_ref_order)
1038+
1039+
def test_update_legacy_self_ref_order_handles_deeper_self_ref_chains(self):
1040+
root = self._make_transferred_store(_self_ref_fk="", _self_ref_order=None)
1041+
child = self._make_transferred_store(
1042+
_self_ref_fk=root.id,
1043+
_self_ref_order=None,
1044+
)
1045+
grandchild = self._make_transferred_store(
1046+
_self_ref_fk=child.id,
1047+
_self_ref_order=None,
1048+
)
1049+
1050+
_update_legacy_self_ref_order(self.transfer_session)
1051+
1052+
root.refresh_from_db()
1053+
child.refresh_from_db()
1054+
grandchild.refresh_from_db()
1055+
self.assertEqual(root._self_ref_order, 0)
1056+
self.assertEqual(child._self_ref_order, 0)
1057+
self.assertEqual(grandchild._self_ref_order, 0)
1058+
10121059
def test_local_dequeue_operation(self):
10131060
self.transfer_session.records_transferred = 1
10141061
self.context.filter = [self.transfer_session.filter]

0 commit comments

Comments
 (0)