Skip to content

Commit b09f956

Browse files
committed
fix: 315764 improve _handle_deduplicated
1 parent aa17171 commit b09f956

2 files changed

Lines changed: 96 additions & 75 deletions

File tree

src/country_workspace/contrib/hope/push/orchestration.py

Lines changed: 28 additions & 18 deletions
Original file line numberDiff line numberDiff line change
@@ -171,11 +171,20 @@ def _lock_and_fail_rdp(rdp_id: int, *, reason: str) -> None:
171171

172172

173173
def _handle_deduplicated(rdp: Rdp, rdp_id: int, findings_count: int) -> None:
174-
origin_job = AsyncJob.objects.filter(rdp_id=rdp_id).order_by("-id").first()
175-
max_findings_percent: int = origin_job.config.get("max_dedup_findings_percent", 0) if origin_job else 0
174+
try:
175+
origin_job = AsyncJob.objects.filter(rdp_id=rdp_id).latest("id")
176+
except AsyncJob.DoesNotExist:
177+
_lock_and_fail_rdp(rdp_id, reason="missing origin AsyncJob for threshold config")
178+
return
179+
180+
max_findings_percent: int = origin_job.config.get("max_dedup_findings_percent", 0)
176181

177182
total_individuals = qs_individuals_for_rdp(rdp=rdp).count()
178-
findings_rate = findings_count / total_individuals * 100 if total_individuals > 0 else float(findings_count)
183+
if total_individuals == 0:
184+
_lock_and_fail_rdp(rdp_id, reason="no individuals linked to RDP; cannot compute findings rate")
185+
return
186+
187+
findings_rate = findings_count / total_individuals * 100
179188

180189
logger.info(
181190
"dedup_callback_handle: rdp_id=%s DEDUPLICATED findings_count=%s total_individuals=%s "
@@ -195,27 +204,28 @@ def _handle_deduplicated(rdp: Rdp, rdp_id: int, findings_count: int) -> None:
195204
return
196205

197206
with transaction.atomic():
198-
locked = lock_rdp_for_update(pk=rdp_id)
199-
if locked.status != Rdp.PushStatus.DEDUP_PENDING:
207+
locked_rdp = lock_rdp_for_update(pk=rdp_id)
208+
if locked_rdp.status != Rdp.PushStatus.DEDUP_PENDING:
200209
logger.warning(
201210
"dedup_callback_handle: rdp_id=%s skip push queue; status changed to %s",
202211
rdp_id,
203-
locked.status,
212+
locked_rdp.status,
204213
)
205214
return
206-
locked.status = Rdp.PushStatus.PENDING
207-
locked.save(update_fields=["status"])
215+
locked_rdp.status = Rdp.PushStatus.PENDING
216+
locked_rdp.save(update_fields=["status"])
217+
218+
push_job = AsyncJob.objects.create(
219+
description=f"Push RDP {rdp_id} to HOPE (post-dedup)",
220+
type=AsyncJob.JobType.TASK,
221+
owner=locked_rdp.pushed_by,
222+
action=fqn(push_existing_rdp_core),
223+
program=locked_rdp.program,
224+
rdp=locked_rdp,
225+
config={"rdp_id": rdp_id},
226+
)
227+
transaction.on_commit(push_job.queue)
208228

209-
push_job = AsyncJob.objects.create(
210-
description=f"Push RDP {rdp_id} to HOPE (post-dedup)",
211-
type=AsyncJob.JobType.TASK,
212-
owner=rdp.pushed_by,
213-
action=fqn(push_existing_rdp_core),
214-
program=rdp.program,
215-
rdp=rdp,
216-
config={"rdp_id": rdp_id},
217-
)
218-
push_job.queue()
219229
logger.info(
220230
"dedup_callback_handle: rdp_id=%s status DEDUP_PENDING→PENDING; queued push job_id=%s",
221231
rdp_id,

tests/contrib/hope/push/test_dedup_callback_flow.py

Lines changed: 68 additions & 57 deletions
Original file line numberDiff line numberDiff line change
@@ -206,6 +206,23 @@ def _patch_dedup_status(mocker: MockerFixture, *, state: str, findings_count: in
206206
return policy
207207

208208

209+
def _patch_origin_job(mocker: MockerFixture, config: dict | None = None, *, missing: bool = False):
210+
qs = mocker.MagicMock()
211+
if missing:
212+
qs.latest.side_effect = AsyncJob.DoesNotExist
213+
else:
214+
qs.latest.return_value = mocker.MagicMock(config=config if config is not None else {})
215+
mocker.patch.object(AsyncJob.objects, "filter", return_value=qs)
216+
return qs
217+
218+
219+
def _patch_individuals_count(mocker: MockerFixture, count: int) -> None:
220+
mocker.patch(
221+
f"{MOD}.qs_individuals_for_rdp",
222+
return_value=mocker.MagicMock(count=mocker.MagicMock(return_value=count)),
223+
)
224+
225+
209226
def test_dedup_callback_handle_rdp_not_found(mocker: MockerFixture) -> None:
210227
_patch_rdp_for_push(mocker, not_found=True)
211228
set_status = mocker.patch(f"{MOD}.set_rdp_push_status")
@@ -320,25 +337,17 @@ def test_dedup_callback_handle_deduplicated_within_threshold_queues_push(mocker:
320337
locked_pending.status = Rdp.PushStatus.DEDUP_PENDING
321338
_patch_rdp_for_push(mocker, rdp)
322339
_patch_dedup_status(mocker, state=DeduplicationSetState.DEDUPLICATED, findings_count=5)
323-
origin_job = mocker.MagicMock(config={"max_dedup_findings_percent": 10})
324-
mocker.patch.object(
325-
AsyncJob.objects,
326-
"filter",
327-
return_value=mocker.MagicMock(
328-
order_by=mocker.MagicMock(return_value=mocker.MagicMock(first=mocker.MagicMock(return_value=origin_job)))
329-
),
330-
)
331-
mocker.patch(
332-
f"{MOD}.qs_individuals_for_rdp", return_value=mocker.MagicMock(count=mocker.MagicMock(return_value=100))
333-
)
340+
_patch_origin_job(mocker, {"max_dedup_findings_percent": 10})
341+
_patch_individuals_count(mocker, 100)
334342
mocker.patch(f"{MOD}.lock_rdp_for_update", return_value=locked_pending)
335343
push_job = mocker.MagicMock()
336344
create_job = mocker.patch.object(AsyncJob.objects, "create", return_value=push_job)
345+
on_commit = mocker.patch(f"{MOD}.transaction.on_commit")
337346

338347
dedup_callback_handle(rdp_id=99)
339348

340349
create_job.assert_called_once()
341-
push_job.queue.assert_called_once()
350+
on_commit.assert_called_once_with(push_job.queue)
342351
assert locked_pending.status == Rdp.PushStatus.PENDING
343352

344353

@@ -350,23 +359,16 @@ def test_dedup_callback_handle_deduplicated_skips_push_when_status_changed_under
350359
locked.status = Rdp.PushStatus.FAILURE
351360
_patch_rdp_for_push(mocker, rdp)
352361
_patch_dedup_status(mocker, state=DeduplicationSetState.DEDUPLICATED, findings_count=5)
353-
origin_job = mocker.MagicMock(config={"max_dedup_findings_percent": 10})
354-
mocker.patch.object(
355-
AsyncJob.objects,
356-
"filter",
357-
return_value=mocker.MagicMock(
358-
order_by=mocker.MagicMock(return_value=mocker.MagicMock(first=mocker.MagicMock(return_value=origin_job)))
359-
),
360-
)
361-
mocker.patch(
362-
f"{MOD}.qs_individuals_for_rdp", return_value=mocker.MagicMock(count=mocker.MagicMock(return_value=100))
363-
)
362+
_patch_origin_job(mocker, {"max_dedup_findings_percent": 10})
363+
_patch_individuals_count(mocker, 100)
364364
mocker.patch(f"{MOD}.lock_rdp_for_update", return_value=locked)
365365
create_job = mocker.patch.object(AsyncJob.objects, "create")
366+
on_commit = mocker.patch(f"{MOD}.transaction.on_commit")
366367

367368
dedup_callback_handle(rdp_id=99)
368369

369370
create_job.assert_not_called()
371+
on_commit.assert_not_called()
370372
locked.save.assert_not_called()
371373

372374

@@ -376,17 +378,8 @@ def test_dedup_callback_handle_deduplicated_exceeds_threshold_marks_failure(mock
376378
locked.status = Rdp.PushStatus.DEDUP_PENDING
377379
_patch_rdp_for_push(mocker, rdp)
378380
_patch_dedup_status(mocker, state=DeduplicationSetState.DEDUPLICATED, findings_count=20)
379-
origin_job = mocker.MagicMock(config={"max_dedup_findings_percent": 10})
380-
mocker.patch.object(
381-
AsyncJob.objects,
382-
"filter",
383-
return_value=mocker.MagicMock(
384-
order_by=mocker.MagicMock(return_value=mocker.MagicMock(first=mocker.MagicMock(return_value=origin_job)))
385-
),
386-
)
387-
mocker.patch(
388-
f"{MOD}.qs_individuals_for_rdp", return_value=mocker.MagicMock(count=mocker.MagicMock(return_value=100))
389-
)
381+
_patch_origin_job(mocker, {"max_dedup_findings_percent": 10})
382+
_patch_individuals_count(mocker, 100)
390383
mocker.patch(f"{MOD}.lock_rdp_for_update", return_value=locked)
391384
set_status = mocker.patch(f"{MOD}.set_rdp_push_status")
392385
create_job = mocker.patch.object(AsyncJob.objects, "create")
@@ -404,17 +397,8 @@ def test_dedup_callback_handle_default_threshold_is_zero(mocker: MockerFixture)
404397
locked.status = Rdp.PushStatus.DEDUP_PENDING
405398
_patch_rdp_for_push(mocker, rdp)
406399
_patch_dedup_status(mocker, state=DeduplicationSetState.DEDUPLICATED, findings_count=1)
407-
origin_job = mocker.MagicMock(config={}) # no max_dedup_findings_percent
408-
mocker.patch.object(
409-
AsyncJob.objects,
410-
"filter",
411-
return_value=mocker.MagicMock(
412-
order_by=mocker.MagicMock(return_value=mocker.MagicMock(first=mocker.MagicMock(return_value=origin_job)))
413-
),
414-
)
415-
mocker.patch(
416-
f"{MOD}.qs_individuals_for_rdp", return_value=mocker.MagicMock(count=mocker.MagicMock(return_value=100))
417-
)
400+
_patch_origin_job(mocker, {})
401+
_patch_individuals_count(mocker, 100)
418402
mocker.patch(f"{MOD}.lock_rdp_for_update", return_value=locked)
419403
set_status = mocker.patch(f"{MOD}.set_rdp_push_status")
420404

@@ -430,21 +414,48 @@ def test_dedup_callback_handle_zero_findings_within_threshold_queues_push(mocker
430414
locked_pending.status = Rdp.PushStatus.DEDUP_PENDING
431415
_patch_rdp_for_push(mocker, rdp)
432416
_patch_dedup_status(mocker, state=DeduplicationSetState.DEDUPLICATED, findings_count=0)
433-
origin_job = mocker.MagicMock(config={}) # default threshold=0
434-
mocker.patch.object(
435-
AsyncJob.objects,
436-
"filter",
437-
return_value=mocker.MagicMock(
438-
order_by=mocker.MagicMock(return_value=mocker.MagicMock(first=mocker.MagicMock(return_value=origin_job)))
439-
),
440-
)
441-
mocker.patch(
442-
f"{MOD}.qs_individuals_for_rdp", return_value=mocker.MagicMock(count=mocker.MagicMock(return_value=100))
443-
)
417+
_patch_origin_job(mocker, {})
418+
_patch_individuals_count(mocker, 100)
444419
mocker.patch(f"{MOD}.lock_rdp_for_update", return_value=locked_pending)
445420
push_job = mocker.MagicMock()
446421
mocker.patch.object(AsyncJob.objects, "create", return_value=push_job)
422+
on_commit = mocker.patch(f"{MOD}.transaction.on_commit")
423+
424+
dedup_callback_handle(rdp_id=99)
425+
426+
on_commit.assert_called_once_with(push_job.queue)
427+
428+
429+
def test_dedup_callback_handle_missing_origin_job_marks_failure(mocker: MockerFixture) -> None:
430+
rdp = _make_rdp(mocker)
431+
locked = mocker.MagicMock()
432+
locked.status = Rdp.PushStatus.DEDUP_PENDING
433+
_patch_rdp_for_push(mocker, rdp)
434+
_patch_dedup_status(mocker, state=DeduplicationSetState.DEDUPLICATED, findings_count=0)
435+
_patch_origin_job(mocker, missing=True)
436+
mocker.patch(f"{MOD}.lock_rdp_for_update", return_value=locked)
437+
set_status = mocker.patch(f"{MOD}.set_rdp_push_status")
438+
create_job = mocker.patch.object(AsyncJob.objects, "create")
439+
440+
dedup_callback_handle(rdp_id=99)
441+
442+
set_status.assert_called_once_with(rdp=locked, status=Rdp.PushStatus.FAILURE, hope_rdi_id="N/A")
443+
create_job.assert_not_called()
444+
445+
446+
def test_dedup_callback_handle_zero_individuals_marks_failure(mocker: MockerFixture) -> None:
447+
rdp = _make_rdp(mocker)
448+
locked = mocker.MagicMock()
449+
locked.status = Rdp.PushStatus.DEDUP_PENDING
450+
_patch_rdp_for_push(mocker, rdp)
451+
_patch_dedup_status(mocker, state=DeduplicationSetState.DEDUPLICATED, findings_count=0)
452+
_patch_origin_job(mocker, {"max_dedup_findings_percent": 10})
453+
_patch_individuals_count(mocker, 0)
454+
mocker.patch(f"{MOD}.lock_rdp_for_update", return_value=locked)
455+
set_status = mocker.patch(f"{MOD}.set_rdp_push_status")
456+
create_job = mocker.patch.object(AsyncJob.objects, "create")
447457

448458
dedup_callback_handle(rdp_id=99)
449459

450-
push_job.queue.assert_called_once()
460+
set_status.assert_called_once_with(rdp=locked, status=Rdp.PushStatus.FAILURE, hope_rdi_id="N/A")
461+
create_job.assert_not_called()

0 commit comments

Comments
 (0)