Skip to content

Commit 867c689

Browse files
committed
feat: Add proper logging for dedup callback
1 parent 4578519 commit 867c689

5 files changed

Lines changed: 250 additions & 21 deletions

File tree

src/country_workspace/config/fragments/constance.py

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
11
from .app import AURORA_API_TOKEN, AURORA_API_URL, HOPE_API_TOKEN, HOPE_API_URL, NEW_USER_DEFAULT_GROUP
2-
from .dedup import DEDUP_API_URL, DEDUP_API_TOKEN
2+
from .dedup import DEDUP_API_URL, DEDUP_API_TOKEN, APP_BASE_URL
33
from .kobo import KOBO_API_TOKEN, KOBO_KF_URL, KOBO_MASTER_API_TOKEN, KOBO_PROJECT_VIEW_ID
44
from .mail import MAILJET_API_KEY, MAILJET_SECRET_KEY
55

@@ -68,6 +68,7 @@
6868
"MAILJET_SECRET_KEY": (MAILJET_SECRET_KEY, "Mailjet secret key", "write_only_text_input"),
6969
"DEDUP_API_URL": (DEDUP_API_URL, "Dedup Engine server address", str),
7070
"DEDUP_API_TOKEN": (DEDUP_API_TOKEN, "Dedup Engine API access token", "write_only_text_input"),
71+
"APP_BASE_URL": (APP_BASE_URL, "Current server address", str),
7172
"CHUNK_SIZE_FOR_VALIDATION_TASK": (500, "Number of records to process per chunk in validation tasks", int),
7273
"RDP_CLEANUP_DAYS": (
7374
30,
@@ -161,6 +162,7 @@
161162
"NEW_USER_IS_STAFF",
162163
"NEW_USER_DEFAULT_GROUP",
163164
"CONCURRENCY_GUARD",
165+
"APP_BASE_URL",
164166
),
165167
}
166168

src/country_workspace/config/urls.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -18,7 +18,7 @@
1818
path("select2/", include(django_select2.urls)),
1919
path(r"__debug__/", include(debug_toolbar.urls)),
2020
path(
21-
"hope/dedup/callback/<str:signed_token>/",
21+
"api/dedup/callback/<str:signed_token>/",
2222
DeduplicationCallbackView.as_view(),
2323
name="dedup_callback",
2424
),

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

Lines changed: 87 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -1,9 +1,10 @@
1+
import logging
12
from collections.abc import Callable, Iterator
23
from functools import partial
34
from typing import Any, NamedTuple
5+
from constance import config
46
from uuid import UUID, uuid4
57

6-
from django.conf import settings
78
from django.core import signing
89
from django.db import IntegrityError, transaction
910
from django.urls import reverse
@@ -43,6 +44,8 @@
4344
_DEDUP_CALLBACK_SIGN_KEY = "dedup_callback"
4445
_DEDUP_CALLBACK_MAX_AGE = 60 * 60 * 96 # 96 hours
4546

47+
logger = logging.getLogger(__name__)
48+
4649

4750
class CloneDeduplicationContext(NamedTuple):
4851
dedup_source_pk: int
@@ -130,7 +133,7 @@ def _build_dedup_callback_url(rdp_id: int, job_id: int) -> str:
130133
key=_DEDUP_CALLBACK_SIGN_KEY,
131134
)
132135
path = reverse("dedup_callback", kwargs={"signed_token": token})
133-
base = getattr(settings, "APP_BASE_URL", "").rstrip("/")
136+
base = config.APP_BASE_URL.rstrip("/")
134137
return f"{base}{path}"
135138

136139

@@ -144,9 +147,10 @@ def create_rdp_and_start_dedup_core(job: AsyncJob) -> dict[str, Any]:
144147
create_result = create_rdp_core(job)
145148
rdp_id = create_result["rdp_id"]
146149

147-
_check, rdp = claim_rdp_deduplication(rdp_id)
148-
if rdp is None:
149-
raise HopePushError({"errors": ["RDP: could not claim deduplication set."], "rdp_id": rdp_id})
150+
with transaction.atomic():
151+
_check, rdp = claim_rdp_deduplication(rdp_id)
152+
if rdp is None:
153+
raise HopePushError({"errors": ["RDP: could not claim deduplication set."], "rdp_id": rdp_id})
150154

151155
job.refresh_from_db()
152156

@@ -179,11 +183,23 @@ def create_and_push_rdp_core(job: AsyncJob) -> dict[str, Any]:
179183
return create_rdp_and_start_dedup_core(job)
180184

181185

182-
def _lock_and_fail_rdp(rdp_id: int) -> None:
186+
def _lock_and_fail_rdp(rdp_id: int, *, reason: str) -> None:
183187
with transaction.atomic():
184188
locked = lock_rdp_for_update(pk=rdp_id)
185189
if locked.status == Rdp.PushStatus.DEDUP_PENDING:
186190
set_rdp_push_status(rdp=locked, status=Rdp.PushStatus.FAILURE, hope_rdi_id="N/A")
191+
logger.warning(
192+
"dedup_callback_handle: rdp_id=%s marked FAILURE (%s)",
193+
rdp_id,
194+
reason,
195+
)
196+
else:
197+
logger.info(
198+
"dedup_callback_handle: rdp_id=%s skip FAILURE (%s); status=%s",
199+
rdp_id,
200+
reason,
201+
locked.status,
202+
)
187203

188204

189205
def _handle_deduplicated(rdp: Rdp, rdp_id: int, findings_count: int) -> None:
@@ -193,13 +209,31 @@ def _handle_deduplicated(rdp: Rdp, rdp_id: int, findings_count: int) -> None:
193209
total_individuals = qs_individuals_for_rdp(rdp=rdp).count()
194210
findings_rate = findings_count / total_individuals * 100 if total_individuals > 0 else float(findings_count)
195211

212+
logger.info(
213+
"dedup_callback_handle: rdp_id=%s DEDUPLICATED findings_count=%s total_individuals=%s "
214+
"findings_rate=%.2f%% max_dedup_findings_percent=%s",
215+
rdp_id,
216+
findings_count,
217+
total_individuals,
218+
findings_rate,
219+
max_findings_percent,
220+
)
221+
196222
if findings_rate > max_findings_percent:
197-
_lock_and_fail_rdp(rdp_id)
223+
_lock_and_fail_rdp(
224+
rdp_id,
225+
reason=f"findings_rate {findings_rate:.2f}% exceeds threshold {max_findings_percent}%",
226+
)
198227
return
199228

200229
with transaction.atomic():
201230
locked = lock_rdp_for_update(pk=rdp_id)
202231
if locked.status != Rdp.PushStatus.DEDUP_PENDING:
232+
logger.warning(
233+
"dedup_callback_handle: rdp_id=%s skip push queue; status changed to %s",
234+
rdp_id,
235+
locked.status,
236+
)
203237
return
204238
locked.status = Rdp.PushStatus.PENDING
205239
locked.save(update_fields=["status"])
@@ -214,6 +248,11 @@ def _handle_deduplicated(rdp: Rdp, rdp_id: int, findings_count: int) -> None:
214248
config={"rdp_id": rdp_id},
215249
)
216250
push_job.queue()
251+
logger.info(
252+
"dedup_callback_handle: rdp_id=%s status DEDUP_PENDING→PENDING; queued push job_id=%s",
253+
rdp_id,
254+
push_job.pk,
255+
)
217256

218257

219258
def dedup_callback_handle(rdp_id: int) -> None:
@@ -231,27 +270,64 @@ def dedup_callback_handle(rdp_id: int) -> None:
231270
try:
232271
rdp = rdp_for_push(pk=rdp_id)
233272
except Rdp.DoesNotExist:
273+
logger.warning("dedup_callback_handle: rdp_id=%s not found", rdp_id)
234274
return
235275

236276
if rdp.status != Rdp.PushStatus.DEDUP_PENDING:
277+
logger.info(
278+
"dedup_callback_handle: rdp_id=%s skip; status=%s (expected DEDUP_PENDING)",
279+
rdp_id,
280+
rdp.status,
281+
)
237282
return
238283

239284
if not rdp.deduplication_set_id:
240-
_lock_and_fail_rdp(rdp_id)
285+
logger.warning("dedup_callback_handle: rdp_id=%s missing deduplication_set_id", rdp_id)
286+
_lock_and_fail_rdp(rdp_id, reason="missing deduplication_set_id")
241287
return
242288

289+
logger.info(
290+
"dedup_callback_handle: rdp_id=%s deduplication_set_id=%s; fetching remote status",
291+
rdp_id,
292+
rdp.deduplication_set_id,
293+
)
294+
243295
status = get_rdp_policy(rdp).deduplication_status(rdp)
244-
if status is None or status.response_status != DedupResponseStatus.OK:
245-
# Dedup engine unavailable — leave in DEDUP_PENDING for the next callback.
296+
if status is None:
297+
logger.warning(
298+
"dedup_callback_handle: rdp_id=%s dedup engine returned no status; leaving DEDUP_PENDING",
299+
rdp_id,
300+
)
301+
return
302+
303+
if status.response_status != DedupResponseStatus.OK:
304+
logger.warning(
305+
"dedup_callback_handle: rdp_id=%s dedup engine status unavailable (%s); leaving DEDUP_PENDING",
306+
rdp_id,
307+
status.response_status,
308+
)
246309
return
247310

248311
dedup_state = status.deduplication_set_status
249312
terminal_failure_states = {DeduplicationSetState.DEDUPLICATION_FAILED, DeduplicationSetState.ENCODING_FAILED}
250313

314+
logger.info(
315+
"dedup_callback_handle: rdp_id=%s remote_state=%s findings_count=%s",
316+
rdp_id,
317+
dedup_state,
318+
status.findings_count,
319+
)
320+
251321
if dedup_state in terminal_failure_states:
252-
_lock_and_fail_rdp(rdp_id)
322+
_lock_and_fail_rdp(rdp_id, reason=f"dedup engine state {dedup_state}")
253323
elif dedup_state == DeduplicationSetState.DEDUPLICATED:
254324
_handle_deduplicated(rdp, rdp_id, status.findings_count)
325+
else:
326+
logger.info(
327+
"dedup_callback_handle: rdp_id=%s intermediate state %s; waiting for next callback",
328+
rdp_id,
329+
dedup_state,
330+
)
255331

256332

257333
def claim_rdp_deduplication(rdp_id: int) -> tuple[ActionCheck, Rdp | None]:

tests/contrib/hope/push/test_dedup_callback_flow.py

Lines changed: 32 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -8,6 +8,7 @@
88
"""
99

1010
import pytest
11+
from constance.test import override_config
1112
from pytest_mock import MockerFixture
1213

1314
from country_workspace.contrib.dedup_engine import (
@@ -35,11 +36,10 @@
3536
# ---------------------------------------------------------------------------
3637

3738

38-
def test_build_dedup_callback_url_contains_signed_token(settings) -> None:
39-
settings.APP_BASE_URL = "https://cw.example.org"
39+
@override_config(APP_BASE_URL="https://cw.example.org")
40+
def test_build_dedup_callback_url_contains_signed_token() -> None:
4041
url = _build_dedup_callback_url(rdp_id=7, job_id=42)
41-
assert url.startswith("https://cw.example.org")
42-
assert "hope/dedup/callback/" in url
42+
assert url.startswith("https://cw.example.org/api/dedup/callback/")
4343

4444

4545
def test_build_dedup_callback_url_token_verifiable() -> None:
@@ -342,6 +342,34 @@ def test_dedup_callback_handle_deduplicated_within_threshold_queues_push(mocker:
342342
assert locked_pending.status == Rdp.PushStatus.PENDING
343343

344344

345+
def test_dedup_callback_handle_deduplicated_skips_push_when_status_changed_under_lock(
346+
mocker: MockerFixture,
347+
) -> None:
348+
rdp = _make_rdp(mocker)
349+
locked = mocker.MagicMock()
350+
locked.status = Rdp.PushStatus.FAILURE
351+
_patch_rdp_for_push(mocker, rdp)
352+
_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+
)
364+
mocker.patch(f"{MOD}.lock_rdp_for_update", return_value=locked)
365+
create_job = mocker.patch.object(AsyncJob.objects, "create")
366+
367+
dedup_callback_handle(rdp_id=99)
368+
369+
create_job.assert_not_called()
370+
locked.save.assert_not_called()
371+
372+
345373
def test_dedup_callback_handle_deduplicated_exceeds_threshold_marks_failure(mocker: MockerFixture) -> None:
346374
rdp = _make_rdp(mocker)
347375
locked = mocker.MagicMock()

0 commit comments

Comments
 (0)