Skip to content

Commit 95df3d3

Browse files
committed
re-implement RefOrderLookUp with batch awareness
1 parent d1325b0 commit 95df3d3

3 files changed

Lines changed: 199 additions & 23 deletions

File tree

.gitignore

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -101,3 +101,5 @@ Pipfile
101101

102102
# Version file generated by setuptools-scm
103103
morango/_version.py
104+
105+
graphify-out/

morango/sync/stream/serialize.py

Lines changed: 110 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
import json
22
import logging
3+
from collections import defaultdict, deque
34
from typing import Generator, Iterable, Iterator, List, Optional, Type
45

56
from django.core.serializers.json import DjangoJSONEncoder
@@ -25,13 +26,22 @@
2526
class SerializeTask(object):
2627
"""Carrier class for providing context through the pipeline"""
2728

28-
__slots__ = ("model", "obj", "store", "counter")
29+
__slots__ = (
30+
"model",
31+
"obj",
32+
"store",
33+
"counter",
34+
"_self_ref_fk_value",
35+
"_self_ref_order",
36+
)
2937

3038
def __init__(self, model: Type[SyncableModel], obj: SyncableModel):
3139
self.model = model
3240
self.obj = obj
3341
self.store: Optional[Store] = None
3442
self.counter: Optional[RecordMaxCounter] = None
43+
self._self_ref_fk_value: Optional[str] = None
44+
self._self_ref_order: Optional[int] = None
3545

3646
@property
3747
def is_store_update(self):
@@ -51,6 +61,20 @@ def self_referential_fk(self) -> Optional[str]:
5161
"""Return the attname of the self-referential FK on *model*, or ``None``."""
5262
return self_referential_fk(self.model)
5363

64+
@property
65+
def self_ref_fk_value(self) -> Optional[str]:
66+
return self._self_ref_fk_value
67+
68+
@property
69+
def self_ref_order(self) -> Optional[int]:
70+
return self._self_ref_order
71+
72+
def set_self_ref_fk_value(self, value: Optional[str]):
73+
self._self_ref_fk_value = value
74+
75+
def set_self_ref_order(self, value: Optional[int]):
76+
self._self_ref_order = value
77+
5478

5579
class AppModelSource(Source[SerializeTask]):
5680
"""
@@ -138,6 +162,84 @@ def transform(self, tasks: List[SerializeTask]) -> List[SerializeTask]:
138162
return tasks
139163

140164

165+
class SelfRefOrderLookup(Transform[List[SerializeTask]]):
166+
"""
167+
Computes self-referential metadata for a buffered batch of tasks.
168+
169+
Resolution order is cache-first, then DB fallback:
170+
- roots (`_self_ref_fk == ""`) get order 0 immediately
171+
- child order is `parent_order + 1` when parent order is known
172+
- unresolved/missing parents keep order as ``None``
173+
"""
174+
175+
def __init__(self):
176+
# Carries resolved parent orders across buffered chunks in the same pipeline run.
177+
self.known_order_by_id = {}
178+
179+
def transform(self, tasks: List[SerializeTask]) -> List[SerializeTask]:
180+
# Cache of resolved order values keyed by record id.
181+
known_order_by_id = self.known_order_by_id
182+
# Self-ref tasks present in the current buffered batch.
183+
tasks_by_id = {}
184+
# Adjacency list for the in-batch self-ref graph: parent_id -> [child tasks].
185+
children_by_parent = defaultdict(list)
186+
187+
# First pass:
188+
# - identify roots (order = 0)
189+
# - build lookup maps:
190+
# tasks_by_id → tasks present in this batch
191+
# children_by_parent → {parent_id → list of child tasks}
192+
for task in tasks:
193+
self_ref_fk = task.self_referential_fk()
194+
if not self_ref_fk:
195+
task.set_self_ref_fk_value(None)
196+
task.set_self_ref_order(None)
197+
continue
198+
199+
self_ref_fk_value = getattr(task.obj, self_ref_fk) or ""
200+
task.set_self_ref_fk_value(self_ref_fk_value)
201+
tasks_by_id[task.obj.id] = task
202+
203+
if not self_ref_fk_value:
204+
task.set_self_ref_order(0)
205+
known_order_by_id[task.obj.id] = 0
206+
continue
207+
208+
209+
task.set_self_ref_order(None)
210+
children_by_parent[self_ref_fk_value].append(task)
211+
212+
# Parent IDs that are referenced by this batch but not present in this batch.
213+
# These must be looked up from existing Store rows.
214+
external_parent_ids = set(children_by_parent.keys()) - set(tasks_by_id.keys())
215+
if external_parent_ids:
216+
for parent_id, parent_order in Store.objects.filter(
217+
id__in=external_parent_ids
218+
).values_list("id", "_self_ref_order"):
219+
if parent_order is not None:
220+
known_order_by_id[parent_id] = parent_order
221+
222+
# Breadth-first propagation from all known parents (roots + DB-resolved parents).
223+
# Each resolved parent unlocks its children as `parent + 1`.
224+
queue = deque(known_order_by_id.keys())
225+
while queue:
226+
parent_id = queue.popleft()
227+
parent_order = known_order_by_id.get(parent_id)
228+
if parent_order is None:
229+
continue
230+
231+
for task in children_by_parent.get(parent_id, []):
232+
if task.self_ref_order is not None:
233+
continue
234+
child_order = parent_order + 1
235+
task.set_self_ref_order(child_order)
236+
child_id = task.obj.id
237+
known_order_by_id[child_id] = child_order
238+
queue.append(child_id)
239+
240+
return tasks
241+
242+
141243
class StoreUpdate(Transform[SerializeTask]):
142244
"""Processes the updates to the Morango store and record counters."""
143245

@@ -188,10 +290,12 @@ def _handle_store_update(self, task: SerializeTask):
188290

189291
self_ref_fk = task.self_referential_fk()
190292
if self_ref_fk:
191-
new_fk_value = getattr(task.obj, self_ref_fk) or ""
293+
new_fk_value = task.self_ref_fk_value
294+
if new_fk_value is None:
295+
new_fk_value = getattr(task.obj, self_ref_fk) or ""
192296
if new_fk_value != task.store._self_ref_fk:
193297
task.store._self_ref_fk = new_fk_value
194-
task.store._self_ref_order = self._compute_self_ref_order(new_fk_value)
298+
task.store._self_ref_order = task.self_ref_order
195299

196300
def _handle_store_create(self, task: SerializeTask):
197301
kwargs = {
@@ -207,29 +311,11 @@ def _handle_store_create(self, task: SerializeTask):
207311

208312
self_ref_fk = task.self_referential_fk()
209313
if self_ref_fk:
210-
self_ref_fk_value = getattr(task.obj, self_ref_fk) or ""
211-
kwargs["_self_ref_fk"] = self_ref_fk_value
212-
kwargs["_self_ref_order"] = self._compute_self_ref_order(self_ref_fk_value)
314+
kwargs["_self_ref_fk"] = task.self_ref_fk_value
315+
kwargs["_self_ref_order"] = task.self_ref_order
213316

214317
task.set_store(Store(**kwargs))
215318

216-
@staticmethod
217-
def _compute_self_ref_order(self_ref_fk_value):
218-
"""
219-
Compute ``_self_ref_order`` for a self-referential store record.
220-
221-
Returns ``0`` when the record has no parent (root), otherwise queries
222-
the parent ``Store`` row and returns the next order value.
223-
"""
224-
if not self_ref_fk_value:
225-
return 0
226-
parent_order = (
227-
Store.objects.filter(id=self_ref_fk_value)
228-
.values_list("_self_ref_order", flat=True)
229-
.first()
230-
)
231-
return parent_order + 1 if parent_order is not None else None
232-
233319

234320
class ModelPartitionBuffer(Buffer[List[SerializeTask]]):
235321
"""Buffers tasks into chunks that have the same model class."""
@@ -419,6 +505,7 @@ def serialize_into_store(
419505
AppModelSource(profile, sync_filter=sync_filter, dirty_only=dirty_only)
420506
.pipe(Buffer(size=500))
421507
.pipe(StoreLookup(current_id))
508+
.pipe(SelfRefOrderLookup())
422509
.pipe(Unbuffer())
423510
.pipe(StoreUpdate(current_id))
424511
.pipe(ModelPartitionBuffer(size=500))

tests/testapp/tests/sync/stream/test_serialize.py

Lines changed: 87 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@
1010
from morango.sync.stream.serialize import (
1111
AppModelSource,
1212
ModelPartitionBuffer,
13+
SelfRefOrderLookup,
1314
SerializeTask,
1415
StoreLookup,
1516
StoreUpdate,
@@ -143,6 +144,82 @@ def test_transform(self, mock_bulk, mock_qs):
143144
self.assertEqual(results[1].counter, None)
144145

145146

147+
class SelfRefOrderLookupTestCase(TestCase):
148+
def _task(self, obj_id=None, parent_id=None):
149+
obj = mock.Mock()
150+
obj.id = obj_id or uuid.uuid4().hex
151+
obj.parent_id = parent_id
152+
return SerializeTask(mock.Mock(), obj)
153+
154+
@mock.patch("morango.sync.stream.serialize.self_referential_fk", return_value="parent_id")
155+
def test_transform__same_batch_parent_child(self, _mock_srf):
156+
parent = self._task(parent_id=None)
157+
child = self._task(parent_id=parent.obj.id)
158+
159+
SelfRefOrderLookup().transform([parent, child])
160+
161+
self.assertEqual(parent.self_ref_fk_value, "")
162+
self.assertEqual(parent.self_ref_order, 0)
163+
self.assertEqual(child.self_ref_fk_value, parent.obj.id)
164+
self.assertEqual(child.self_ref_order, 1)
165+
166+
@mock.patch("morango.sync.stream.serialize.self_referential_fk", return_value="parent_id")
167+
def test_transform__same_batch_deeper_chain(self, _mock_srf):
168+
root = self._task(parent_id=None)
169+
child = self._task(parent_id=root.obj.id)
170+
grandchild = self._task(parent_id=child.obj.id)
171+
172+
SelfRefOrderLookup().transform([root, child, grandchild])
173+
174+
self.assertEqual(root.self_ref_order, 0)
175+
self.assertEqual(child.self_ref_order, 1)
176+
self.assertEqual(grandchild.self_ref_order, 2)
177+
178+
@mock.patch("morango.sync.stream.serialize.self_referential_fk", return_value="parent_id")
179+
def test_transform__previous_batches_feed_later_children(self, _mock_srf):
180+
root = self._task(parent_id=None)
181+
child = self._task(parent_id=root.obj.id)
182+
grandchild = self._task(parent_id=child.obj.id)
183+
lookup = SelfRefOrderLookup()
184+
185+
lookup.transform([root])
186+
lookup.transform([child])
187+
lookup.transform([grandchild])
188+
189+
self.assertEqual(root.self_ref_order, 0)
190+
self.assertEqual(child.self_ref_order, 1)
191+
self.assertEqual(grandchild.self_ref_order, 2)
192+
193+
@mock.patch("morango.sync.stream.serialize.self_referential_fk", return_value="parent_id")
194+
def test_transform__parent_in_store(self, _mock_srf):
195+
parent_store = _make_store(_self_ref_order=3)
196+
child = self._task(parent_id=parent_store.id)
197+
198+
SelfRefOrderLookup().transform([child])
199+
200+
self.assertEqual(child.self_ref_fk_value, parent_store.id)
201+
self.assertEqual(child.self_ref_order, 4)
202+
203+
@mock.patch("morango.sync.stream.serialize.self_referential_fk", return_value="parent_id")
204+
def test_transform__missing_parent(self, _mock_srf):
205+
missing_parent_id = uuid.uuid4().hex
206+
child = self._task(parent_id=missing_parent_id)
207+
208+
SelfRefOrderLookup().transform([child])
209+
210+
self.assertEqual(child.self_ref_fk_value, missing_parent_id)
211+
self.assertIsNone(child.self_ref_order)
212+
213+
@mock.patch("morango.sync.stream.serialize.self_referential_fk", return_value=None)
214+
def test_transform__non_self_ref(self, _mock_srf):
215+
task = self._task()
216+
217+
SelfRefOrderLookup().transform([task])
218+
219+
self.assertIsNone(task.self_ref_fk_value)
220+
self.assertIsNone(task.self_ref_order)
221+
222+
146223
class StoreUpdateTestCase(SimpleTestCase):
147224
def _build_sync_obj(self, **overrides):
148225
obj = mock.Mock()
@@ -227,6 +304,8 @@ def test_handle_store_create__self_ref_no_parent(self, _mock_srf):
227304
obj = self._build_sync_obj(parent_id=None) # no parent
228305

229306
task = SerializeTask(mock.Mock(), obj)
307+
task.set_self_ref_fk_value("")
308+
task.set_self_ref_order(0)
230309
update._handle_store_create(task)
231310

232311
self.assertEqual(task.store._self_ref_fk, "")
@@ -247,6 +326,8 @@ def test_handle_store_update__self_ref_fk_unchanged(self, _mock_srf):
247326
obj = self._build_sync_obj(parent_id=parent_id) # same FK — no change
248327

249328
task = SerializeTask(mock.Mock(), obj)
329+
task.set_self_ref_fk_value(parent_id)
330+
task.set_self_ref_order(5)
250331
task.set_store(store)
251332
update._handle_store_update(task)
252333

@@ -290,6 +371,8 @@ def test_handle_store_create__self_ref_with_parent(self, _mock_srf):
290371
parent_store = _make_store(_self_ref_order=3)
291372

292373
task = self._task("parent_id", parent_store.id)
374+
task.set_self_ref_fk_value(parent_store.id)
375+
task.set_self_ref_order(4)
293376
self.update._handle_store_create(task)
294377

295378
self.assertEqual(task.store._self_ref_fk, parent_store.id)
@@ -300,6 +383,8 @@ def test_handle_store_create__self_ref_parent_not_in_store(self, _mock_srf):
300383
missing_parent_id = uuid.uuid4().hex
301384

302385
task = self._task("parent_id", missing_parent_id)
386+
task.set_self_ref_fk_value(missing_parent_id)
387+
task.set_self_ref_order(None)
303388
self.update._handle_store_create(task)
304389

305390
self.assertEqual(task.store._self_ref_fk, missing_parent_id)
@@ -321,6 +406,8 @@ def test_handle_store_update__self_ref_fk_changed(self, _mock_srf):
321406
obj.parent_id = new_parent_store.id # FK changed
322407

323408
task = SerializeTask(mock.Mock(), obj)
409+
task.set_self_ref_fk_value(new_parent_store.id)
410+
task.set_self_ref_order(8)
324411
task.set_store(store)
325412
self.update._handle_store_update(task)
326413

0 commit comments

Comments
 (0)