Skip to content

Commit b46528c

Browse files
authored
feat(M2-9862): Out-of-Order Flow Submissions with Duplicate Activities (#1954)
* feat(M2-9862): Fix out-of-order flow submissions and handle duplicate activities - Added FlowSubmissionProgress helper for backend-owned flow tracking - Replaced order-based logic with Counter-based validation - Persisted backend-computed flow completion - Logged mismatches for legacy clients - Added new out-of-order tests and updated existing ones * fix: add back create_at and convert Unix timestamp to datetime in answer existence check * M2-9862 make the create_at optional * fix(M2-9862): fix linting using ruff
1 parent 792344a commit b46528c

8 files changed

Lines changed: 394 additions & 105 deletions

File tree

src/apps/answers/api.py

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -873,12 +873,12 @@ async def answers_existence_check(
873873
await AppletService(session, user.id).exist_by_id(schema.applet_id)
874874
await CheckAccessService(session, user.id).check_answer_check_access(schema.applet_id)
875875
is_exist = await AnswerService(session, user.id, answer_session).is_answers_uploaded(
876-
schema.applet_id, schema.activity_id, schema.submit_id
876+
schema.applet_id, schema.activity_id, schema.submit_id, schema.created_at
877877
)
878878

879879
logger.info(
880880
f"check-existence: applet_id={schema.applet_id}, activity_id={schema.activity_id}, user_id={user.id}, "
881-
f"submit_id={schema.submit_id}, exists={is_exist}, ip={client_ip}"
881+
f"submit_id={schema.submit_id}, created_at={schema.created_at}, exists={is_exist}, ip={client_ip}"
882882
)
883883

884884
return Response[AnswerExistenceResponse](result=AnswerExistenceResponse(exists=is_exist))

src/apps/answers/crud/answers.py

Lines changed: 8 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -518,16 +518,22 @@ async def get_by_applet_activity_submit_or_user_id(
518518
activity_id: str,
519519
user_id: uuid.UUID | None = None,
520520
submit_id: uuid.UUID | None = None,
521+
created_at: int | None = None,
521522
) -> list[AnswerSchema]:
522-
# We're not using created_at for filtering as it causes issues with mobile submissions
523-
# The combination of applet_id, activity_id, and either user_id or submit_id should be sufficient
523+
# created_at is used to distinguish between duplicate activities in flows
524524
query: Query = select(AnswerSchema)
525525
query = query.where(AnswerSchema.applet_id == applet_id)
526526
query = query.filter(AnswerSchema.activity_history_id.startswith(activity_id))
527527
if submit_id:
528528
query = query.where(AnswerSchema.submit_id == submit_id)
529529
if user_id:
530530
query = query.where(AnswerSchema.respondent_id == user_id)
531+
if created_at is not None:
532+
# Convert Unix timestamp (milliseconds) to datetime for comparison
533+
created_at_datetime = datetime.datetime.fromtimestamp(
534+
created_at / 1000.0, tz=datetime.timezone.utc
535+
).replace(tzinfo=None)
536+
query = query.where(AnswerSchema.created_at == created_at_datetime)
531537

532538
db_result = await self._execute(query)
533539
return db_result.scalars().all()

src/apps/answers/domain/answers.py

Lines changed: 1 addition & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -673,9 +673,7 @@ class AppletCompletedEntities(InternalModel):
673673

674674
class AnswersCheck(PublicModel):
675675
applet_id: uuid.UUID
676-
# TODO: created_at can be safely removed after
677-
# the corresponding mobile PR is merged
678-
# https://mindlogger.atlassian.net/browse/M2-9693
676+
# Used to distinguish between duplicate activities in flows
679677
created_at: int | None = None
680678
activity_id: str
681679
submit_id: uuid.UUID | None = None
Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,74 @@
1+
from __future__ import annotations
2+
3+
from collections import Counter
4+
from typing import Iterable
5+
6+
from sqlalchemy.ext.asyncio import AsyncSession
7+
8+
from apps.activity_flows.crud import FlowsHistoryCRUD
9+
from apps.answers.crud.answers import AnswersCRUD
10+
11+
12+
class FlowSubmissionProgress:
13+
"""Track flow submission progress using occurrence counting."""
14+
15+
def __init__(self, session: AsyncSession, answer_session: AsyncSession):
16+
self._session = session
17+
self._answer_session = answer_session
18+
self._expected_counts: Counter[str] = Counter()
19+
self._expected_ids: set[str] = set()
20+
self._submitted_counts: Counter[str] = Counter()
21+
self._has_flow_history: bool = False
22+
23+
async def load(self, flow_history_id: str, submit_id) -> None:
24+
"""Load flow structure and existing submissions for the submit id."""
25+
flow_histories = await FlowsHistoryCRUD(self._session).load_full([flow_history_id], load_activities=False)
26+
if not flow_histories:
27+
self._has_flow_history = False
28+
return
29+
30+
flow_history = flow_histories[0]
31+
self._expected_counts = Counter(item.activity_id for item in flow_history.items)
32+
self._expected_ids = set(self._expected_counts.keys())
33+
self._has_flow_history = True
34+
35+
existing_answers = await AnswersCRUD(self._answer_session).get_by_submit_id(submit_id)
36+
self._submitted_counts = Counter(answer.activity_history_id for answer in existing_answers or [])
37+
38+
@property
39+
def has_flow_history(self) -> bool:
40+
return self._has_flow_history
41+
42+
def is_complete_before_current(self) -> bool:
43+
return self._all_expected_satisfied(self._submitted_counts.items())
44+
45+
def can_accept(self, activity_history_id: str) -> bool:
46+
expected_total = self._expected_counts.get(activity_history_id, 0)
47+
if expected_total == 0:
48+
return False
49+
return self._submitted_counts.get(activity_history_id, 0) < expected_total
50+
51+
def completion_state_after_add(self, activity_history_id: str) -> bool:
52+
temp_counts = self._submitted_counts.copy()
53+
temp_counts[activity_history_id] += 1
54+
return self._all_expected_satisfied(temp_counts.items())
55+
56+
def contains_activity(self, activity_history_id: str) -> bool:
57+
if not self._has_flow_history:
58+
return False
59+
return activity_history_id in self._expected_ids
60+
61+
@property
62+
def expected_total(self) -> int:
63+
return sum(self._expected_counts.values())
64+
65+
@property
66+
def submitted_total(self) -> int:
67+
return sum(self._submitted_counts.values())
68+
69+
def _all_expected_satisfied(self, submitted_items: Iterable[tuple[str, int]]) -> bool:
70+
submitted_map = dict(submitted_items)
71+
for activity_history_id, expected_count in self._expected_counts.items():
72+
if submitted_map.get(activity_history_id, 0) < expected_count:
73+
return False
74+
return True

src/apps/answers/service.py

Lines changed: 80 additions & 73 deletions
Original file line numberDiff line numberDiff line change
@@ -6,7 +6,7 @@
66
import os
77
import time
88
import uuid
9-
from collections import Counter, defaultdict
9+
from collections import defaultdict
1010
from json import JSONDecodeError
1111
from typing import Callable, List, Mapping, Optional
1212

@@ -89,6 +89,7 @@
8989
WrongRespondentForAnswerGroup,
9090
)
9191
from apps.answers.filters import AppletSubmitDateFilter, ReviewAppletItemFilter, SummaryActivityFilter
92+
from apps.answers.flow_submission_progress import FlowSubmissionProgress
9293
from apps.answers.tasks import create_report
9394
from apps.applets.crud import AppletsCRUD
9495
from apps.applets.domain.applet_history import Version
@@ -150,46 +151,54 @@ async def create_answer(self, activity_answer: AppletAnswerCreate, device_id: st
150151
async def _create_respondent_answer(
151152
self, activity_answer: AppletAnswerCreate, device_id: str | None
152153
) -> AnswerSchema:
153-
await self._validate_respondent_answer(activity_answer)
154-
return await self._create_answer(activity_answer, device_id)
154+
flow_progress = await self._validate_respondent_answer(activity_answer)
155+
return await self._create_answer(activity_answer, device_id, flow_progress)
155156

156157
async def _create_anonymous_answer(
157158
self, activity_answer: AppletAnswerCreate, device_id: str | None
158159
) -> AnswerSchema:
159-
await self._validate_anonymous_answer(activity_answer)
160-
return await self._create_answer(activity_answer, device_id)
160+
flow_progress = await self._validate_anonymous_answer(activity_answer)
161+
return await self._create_answer(activity_answer, device_id, flow_progress)
161162

162-
async def _validate_respondent_answer(self, activity_answer: AppletAnswerCreate) -> None:
163-
await self._validate_answer(activity_answer)
163+
async def _validate_respondent_answer(self, activity_answer: AppletAnswerCreate) -> FlowSubmissionProgress | None:
164+
flow_progress = await self._validate_answer(activity_answer)
164165
await self._validate_applet_for_user_response(activity_answer.applet_id)
165166

166-
async def _validate_anonymous_answer(self, activity_answer: AppletAnswerCreate) -> None:
167+
return flow_progress
168+
169+
async def _validate_anonymous_answer(self, activity_answer: AppletAnswerCreate) -> FlowSubmissionProgress | None:
167170
await self._validate_applet_for_anonymous_response(activity_answer.applet_id, activity_answer.version)
168-
await self._validate_answer(activity_answer)
171+
return await self._validate_answer(activity_answer)
169172

170-
async def _validate_answer(self, applet_answer: AppletAnswerCreate) -> None: # noqa: C901
173+
async def _validate_answer(self, applet_answer: AppletAnswerCreate) -> FlowSubmissionProgress | None: # noqa: C901
171174
pk = self._generate_history_id(applet_answer.version)
175+
176+
# Timestamp-based duplicate detection: when created_at provided, check for exact duplicate
177+
# Each unique timestamp represents a distinct submission, bypassing occurrence limits
178+
if applet_answer.created_at is not None:
179+
created_at_ms = int(applet_answer.created_at.timestamp() * 1000)
180+
existing_with_timestamp = await AnswersCRUD(self.answer_session).get_by_applet_activity_submit_or_user_id(
181+
applet_answer.applet_id,
182+
str(applet_answer.activity_id),
183+
None,
184+
applet_answer.submit_id,
185+
created_at_ms,
186+
)
187+
if existing_with_timestamp:
188+
raise ValidationError("Duplicate answer with same timestamp already exists")
189+
172190
existed_answers = await AnswersCRUD(self.answer_session).get_by_submit_id(applet_answer.submit_id)
173191

174192
activity_history_id = pk(applet_answer.activity_id)
175193
flow_history_id = pk(applet_answer.flow_id) if applet_answer.flow_id else None
176-
177-
activity_indexes = set() # same activity is allowed multiple times in flow
178-
latest_activity_index = None
179-
if flow_history_id:
180-
flow_histories = await FlowsHistoryCRUD(self.session).load_full(
181-
[pk(applet_answer.flow_id)], load_activities=False
182-
)
183-
if not flow_histories:
194+
flow_progress: FlowSubmissionProgress | None = None
195+
# Only use occurrence-based flow validation when created_at is NOT provided
196+
if flow_history_id and applet_answer.created_at is None:
197+
flow_progress = FlowSubmissionProgress(self.session, self.answer_session)
198+
await flow_progress.load(flow_history_id, applet_answer.submit_id)
199+
if not flow_progress.has_flow_history:
184200
raise ValidationError("Flow not found")
185-
flow_history = next(iter(flow_histories))
186-
187-
# check activity in the flow
188-
for i, item in enumerate(flow_history.items):
189-
if item.activity_id == activity_history_id:
190-
activity_indexes.add(i)
191-
latest_activity_index = len(flow_history.items) - 1
192-
if not activity_indexes:
201+
if not flow_progress.contains_activity(activity_history_id):
193202
raise ValidationError("Activity not found in the flow")
194203

195204
if existed_answers:
@@ -208,60 +217,28 @@ async def _validate_answer(self, applet_answer: AppletAnswerCreate) -> None: #
208217
if flow_history_id != existed_answer.flow_history_id:
209218
raise ValidationError("Submit id duplicate error")
210219

211-
# check current answer is provided in right order in the flow, so prev activities already answered
212-
prev_answers_count = len(existed_answers)
213-
is_flow_completed = any(answer.is_flow_completed for answer in existed_answers)
214-
215-
# Smart flow completion check
216-
if is_flow_completed:
217-
# Count expected and already persisted activity occurrences
218-
flow_activity_counts = Counter(item.activity_id for item in flow_history.items)
219-
submitted_counts = Counter(answer.activity_history_id for answer in existed_answers)
220-
221-
# If all activities already persisted before this submission, the flow truly finished
222-
if all(submitted_counts.get(act_id, 0) >= count for act_id, count in flow_activity_counts.items()):
223-
raise ValidationError("Flow is already completed")
220+
if flow_progress:
221+
if flow_progress.is_complete_before_current():
222+
raise ValidationError("Flow is already completed")
224223

225-
current_expected_total = flow_activity_counts.get(activity_history_id, 0)
226-
current_submitted = submitted_counts.get(activity_history_id, 0)
227-
228-
# Reject duplicates that exceed expected occurrences for the activity
229-
if current_submitted >= current_expected_total:
230-
raise ValidationError("Flow is already completed")
224+
if not flow_progress.can_accept(activity_history_id):
225+
raise ValidationError("Activity submission exceeds expected occurrences for this flow")
231226

227+
if existed_answers:
232228
logger.info(
233-
"Allowing late submission for flow %s, activity %s, submit_id %s",
229+
"Allowing flow submission for flow %s, activity %s, submit_id %s",
234230
flow_history_id,
235231
activity_history_id,
236232
applet_answer.submit_id,
237233
)
238234

239-
# Continue with existing order validation only if flow not marked complete
240-
# When is_flow_completed=True, we allow out-of-order submissions
241-
if not is_flow_completed and prev_answers_count not in activity_indexes:
242-
assert latest_activity_index is not None
243-
# allow latest activity for flow autocompletion FE logic
244-
if not (
245-
prev_answers_count < latest_activity_index + 1
246-
and max(activity_indexes) == latest_activity_index
247-
and applet_answer.is_flow_completed
248-
):
249-
raise ValidationError("Wrong activity order in the flow")
250-
251-
elif flow_history_id and 0 not in activity_indexes:
252-
# check first flow answer - but allow if flow is marked as complete
253-
# Check both existing answers and current answer for is_flow_completed
254-
if not (
255-
(existed_answers and any(answer.is_flow_completed for answer in existed_answers))
256-
or applet_answer.is_flow_completed
257-
):
258-
raise ValidationError("Wrong activity order in the flow")
259-
260235
activity_history = await ActivityHistoriesCRUD(self.session).get_by_id(activity_history_id)
261236

262237
if not activity_history.applet_id.startswith(f"{applet_answer.applet_id}"):
263238
raise ActivityHistoryDoeNotExist()
264239

240+
return flow_progress
241+
265242
async def _validate_applet_for_anonymous_response(self, applet_id: uuid.UUID, version: str) -> None:
266243
await AppletHistoryService(self.session, applet_id, version).get()
267244
# Validate applet for anonymous answer
@@ -346,12 +323,41 @@ async def _get_answer_relation(
346323

347324
return relation.relation
348325

349-
async def _create_answer(self, applet_answer: AppletAnswerCreate, device_id: str | None) -> AnswerSchema:
326+
async def _create_answer(
327+
self,
328+
applet_answer: AppletAnswerCreate,
329+
device_id: str | None,
330+
flow_progress: FlowSubmissionProgress | None,
331+
) -> AnswerSchema:
350332
assert self.user_id
351333
pk = self._generate_history_id(applet_answer.version)
352334
created_at = applet_answer.created_at or datetime.datetime.now(datetime.UTC).replace(tzinfo=None)
353335
subject_crud = SubjectsCrud(self.session)
354336

337+
activity_history_id = pk(applet_answer.activity_id)
338+
flow_history_id = pk(applet_answer.flow_id) if applet_answer.flow_id else None
339+
client_flow_completed_flag = bool(applet_answer.is_flow_completed) if applet_answer.flow_id else None
340+
341+
is_flow_completed_backend = None
342+
if applet_answer.flow_id:
343+
if flow_progress:
344+
is_flow_completed_backend = flow_progress.completion_state_after_add(activity_history_id)
345+
if client_flow_completed_flag is not None and client_flow_completed_flag != is_flow_completed_backend:
346+
logger.info(
347+
"Flow completion mismatch for flow %s, activity %s, submit_id %s: client=%s backend=%s",
348+
flow_history_id,
349+
activity_history_id,
350+
applet_answer.submit_id,
351+
client_flow_completed_flag,
352+
is_flow_completed_backend,
353+
)
354+
else:
355+
is_flow_completed_backend = client_flow_completed_flag
356+
357+
migrated_data = None
358+
if client_flow_completed_flag is not None:
359+
migrated_data = {"client_flow_completed_flag": client_flow_completed_flag}
360+
355361
respondent_subject = await subject_crud.get_user_subject(
356362
user_id=self.user_id, applet_id=applet_answer.applet_id
357363
)
@@ -413,18 +419,19 @@ async def _create_answer(self, applet_answer: AppletAnswerCreate, device_id: str
413419
applet_id=applet_answer.applet_id,
414420
version=applet_answer.version,
415421
applet_history_id=pk(applet_answer.applet_id),
416-
flow_history_id=pk(applet_answer.flow_id) if applet_answer.flow_id else None,
417-
activity_history_id=pk(applet_answer.activity_id),
422+
flow_history_id=flow_history_id,
423+
activity_history_id=activity_history_id,
418424
respondent_id=self.user_id,
419425
client=applet_answer.client.dict(),
420-
is_flow_completed=bool(applet_answer.is_flow_completed) if applet_answer.flow_id else None,
426+
is_flow_completed=is_flow_completed_backend,
421427
target_subject_id=target_subject.id,
422428
source_subject_id=source_subject.id,
423429
input_subject_id=input_subject.id,
424430
relation=relation,
425431
consent_to_share=applet_answer.consent_to_share,
426432
event_history_id=applet_answer.event_history_id,
427433
device_id=device_id,
434+
migrated_data=migrated_data,
428435
)
429436
)
430437
item_answer = applet_answer.answer
@@ -1724,11 +1731,11 @@ async def get_completed_answers_data_list(
17241731
return result
17251732

17261733
async def is_answers_uploaded(
1727-
self, applet_id: uuid.UUID, activity_id: str, submit_id: uuid.UUID | None = None
1734+
self, applet_id: uuid.UUID, activity_id: str, submit_id: uuid.UUID | None = None, created_at: int | None = None
17281735
) -> bool:
17291736
# check by submit id if provided otherwise by user_id
17301737
answers = await AnswersCRUD(self.answer_session).get_by_applet_activity_submit_or_user_id(
1731-
applet_id, activity_id, self.user_id if not submit_id else None, submit_id
1738+
applet_id, activity_id, self.user_id if not submit_id else None, submit_id, created_at
17321739
)
17331740
if not answers:
17341741
return False

src/apps/answers/tests/conftest.py

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -235,7 +235,7 @@ def answer_create(
235235
submit_id=uuid.uuid4(),
236236
activity_id=applet.activities[0].id,
237237
answer=answer_item_create,
238-
created_at=datetime.datetime.now(datetime.UTC).replace(microsecond=0),
238+
created_at=None, # None by default - uses occurrence-based validation
239239
client=client_meta,
240240
consent_to_share=False,
241241
)

0 commit comments

Comments
 (0)