Skip to content

Commit 40be75d

Browse files
committed
implement server compute order during the dequeue stage
1 parent 6397c61 commit 40be75d

5 files changed

Lines changed: 90 additions & 2 deletions

File tree

morango/constants/capabilities.py

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,3 +2,4 @@
22
ALLOW_CERTIFICATE_PUSHING = "ALLOW_CERTIFICATE_PUSHING"
33
ASYNC_OPERATIONS = "ASYNC_OPERATIONS"
44
FSIC_V2_FORMAT = "FSIC_V2_FORMAT"
5+
SELF_REF_ORDER = "SELF_REF_ORDER"

morango/sync/operations.py

Lines changed: 50 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -19,6 +19,7 @@
1919
from morango.constants import transfer_statuses
2020
from morango.constants.capabilities import ASYNC_OPERATIONS
2121
from morango.constants.capabilities import FSIC_V2_FORMAT
22+
from morango.constants.capabilities import SELF_REF_ORDER
2223
from morango.errors import MorangoInvalidFSICPartition
2324
from morango.errors import MorangoLimitExceeded
2425
from morango.errors import MorangoResumeSyncError
@@ -679,7 +680,52 @@ def _queue_into_buffer_v2(transfersession, chunk_size=200):
679680
)
680681

681682

682-
def _dequeue_into_store(transfer_session, fsic, v2_format=False):
683+
def _update_legacy_self_ref_order(transfer_session):
684+
transferred_store_records = Store.objects.filter(
685+
last_transfer_session_id=transfer_session.id
686+
)
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
695+
696+
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
726+
727+
728+
def _dequeue_into_store(transfer_session, fsic, v2_format=False, self_ref_order=True):
683729
"""
684730
Takes data from the buffers and merges into the store and record max counters.
685731
@@ -703,6 +749,8 @@ def _dequeue_into_store(transfer_session, fsic, v2_format=False):
703749
DBBackend._dequeuing_delete_mc_buffer(cursor, transfer_session.id)
704750
DBBackend._dequeuing_insert_remaining_buffer(cursor, transfer_session.id)
705751
DBBackend._dequeuing_insert_remaining_rmcb(cursor, transfer_session.id)
752+
if not self_ref_order:
753+
_update_legacy_self_ref_order(transfer_session)
706754
DBBackend._dequeuing_delete_remaining_rmcb(cursor, transfer_session.id)
707755
DBBackend._dequeuing_delete_remaining_buffer(cursor, transfer_session.id)
708756

@@ -1041,6 +1089,7 @@ def handle(self, context):
10411089
context.transfer_session,
10421090
fsic,
10431091
v2_format=FSIC_V2_FORMAT in context.capabilities,
1092+
self_ref_order=SELF_REF_ORDER in context.capabilities,
10441093
)
10451094

10461095
return transfer_statuses.COMPLETED

morango/utils.py

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
ASYNC_OPERATIONS,
1010
FSIC_V2_FORMAT,
1111
GZIP_BUFFER_POST,
12+
SELF_REF_ORDER,
1213
)
1314

1415

@@ -61,6 +62,8 @@ def get_capabilities():
6162
if not SETTINGS.MORANGO_DISALLOW_ASYNC_OPERATIONS:
6263
capabilities.add(ASYNC_OPERATIONS)
6364

65+
capabilities.add(SELF_REF_ORDER)
66+
6467
return capabilities
6568

6669

tests/testapp/tests/sync/test_operations.py

Lines changed: 32 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@
88
from facility_profile.models import ConditionalLog, Facility, MyUser, SummaryLog
99

1010
from morango.constants import transfer_statuses
11-
from morango.constants.capabilities import FSIC_V2_FORMAT
11+
from morango.constants.capabilities import FSIC_V2_FORMAT, SELF_REF_ORDER
1212
from morango.errors import MorangoLimitExceeded
1313
from morango.models.certificates import Filter
1414
from morango.models.core import (
@@ -991,6 +991,24 @@ def test_dequeue_into_store(self):
991991
).exists()
992992
)
993993

994+
def test_dequeue_into_store__self_ref_order_fallback_for_missing_capability(self):
995+
Buffer.objects.filter(model_uuid=self.data["model3"]).update(
996+
_self_ref_fk="", _self_ref_order=None
997+
)
998+
Buffer.objects.filter(model_uuid=self.data["model4"]).update(
999+
_self_ref_fk=self.data["model3"], _self_ref_order=None
1000+
)
1001+
1002+
_dequeue_into_store(
1003+
self.transfer_session,
1004+
self.transfer_session.client_fsic,
1005+
v2_format=False,
1006+
self_ref_order=False,
1007+
)
1008+
1009+
self.assertEqual(Store.objects.get(id=self.data["model3"])._self_ref_order, 0)
1010+
self.assertEqual(Store.objects.get(id=self.data["model4"])._self_ref_order, 0)
1011+
9941012
def test_local_dequeue_operation(self):
9951013
self.transfer_session.records_transferred = 1
9961014
self.context.filter = [self.transfer_session.filter]
@@ -1000,6 +1018,19 @@ def test_local_dequeue_operation(self):
10001018
Buffer.objects.filter(transfer_session_id=self.transfer_session.id).exists()
10011019
)
10021020

1021+
@mock.patch("morango.sync.operations._dequeue_into_store")
1022+
def test_local_dequeue_operation__passes_self_ref_order_capability(self, mock_dequeue):
1023+
self.transfer_session.records_transferred = 1
1024+
self.context.capabilities = {SELF_REF_ORDER}
1025+
operation = ReceiverDequeueOperation()
1026+
self.assertEqual(transfer_statuses.COMPLETED, operation.handle(self.context))
1027+
mock_dequeue.assert_called_once_with(
1028+
self.transfer_session,
1029+
self.transfer_session.client_fsic,
1030+
v2_format=False,
1031+
self_ref_order=True,
1032+
)
1033+
10031034
@mock.patch("morango.sync.operations._dequeue_into_store")
10041035
def test_local_dequeue_operation__noop(self, mock_dequeue):
10051036
self.context.is_server = False

tests/testapp/tests/test_utils.py

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@
1212
ALLOW_CERTIFICATE_PUSHING,
1313
ASYNC_OPERATIONS,
1414
FSIC_V2_FORMAT,
15+
SELF_REF_ORDER,
1516
)
1617
from morango.utils import (
1718
CAPABILITIES_CLIENT_HEADER,
@@ -73,6 +74,9 @@ def test_get_capabilities__fsic_v2_format(self):
7374
with self.settings(MORANGO_DISABLE_FSIC_V2_FORMAT=True):
7475
self.assertNotIn(FSIC_V2_FORMAT, get_capabilities())
7576

77+
def test_get_capabilities__self_ref_order(self):
78+
self.assertIn(SELF_REF_ORDER, get_capabilities())
79+
7680
@mock.patch("morango.utils.CAPABILITIES", ("TEST", "SERIALIZE"))
7781
def test_serialize(self):
7882
req = Request()

0 commit comments

Comments
 (0)