From 7ae4b5adab944d89b4af48078fefcdc3984a55d9 Mon Sep 17 00:00:00 2001 From: Mohamed Elmi Date: Tue, 7 Oct 2025 23:54:03 -0500 Subject: [PATCH 1/4] 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 --- src/apps/answers/flow_submission_progress.py | 74 ++++++++++ src/apps/answers/service.py | 131 ++++++++---------- src/apps/answers/tests/test_answers.py | 42 ++++-- .../answers/tests/test_flow_out_of_order.py | 58 ++++++-- 4 files changed, 215 insertions(+), 90 deletions(-) create mode 100644 src/apps/answers/flow_submission_progress.py diff --git a/src/apps/answers/flow_submission_progress.py b/src/apps/answers/flow_submission_progress.py new file mode 100644 index 000000000000..09eb5784db1a --- /dev/null +++ b/src/apps/answers/flow_submission_progress.py @@ -0,0 +1,74 @@ +from __future__ import annotations + +from collections import Counter +from typing import Iterable + +from sqlalchemy.ext.asyncio import AsyncSession + +from apps.activity_flows.crud import FlowsHistoryCRUD +from apps.answers.crud.answers import AnswersCRUD + + +class FlowSubmissionProgress: + """Track flow submission progress using occurrence counting.""" + + def __init__(self, session: AsyncSession, answer_session: AsyncSession): + self._session = session + self._answer_session = answer_session + self._expected_counts: Counter[str] = Counter() + self._expected_ids: set[str] = set() + self._submitted_counts: Counter[str] = Counter() + self._has_flow_history: bool = False + + async def load(self, flow_history_id: str, submit_id) -> None: + """Load flow structure and existing submissions for the submit id.""" + flow_histories = await FlowsHistoryCRUD(self._session).load_full([flow_history_id], load_activities=False) + if not flow_histories: + self._has_flow_history = False + return + + flow_history = flow_histories[0] + self._expected_counts = Counter(item.activity_id for item in flow_history.items) + self._expected_ids = set(self._expected_counts.keys()) + self._has_flow_history = True + + existing_answers = await AnswersCRUD(self._answer_session).get_by_submit_id(submit_id) + self._submitted_counts = Counter(answer.activity_history_id for answer in existing_answers or []) + + @property + def has_flow_history(self) -> bool: + return self._has_flow_history + + def is_complete_before_current(self) -> bool: + return self._all_expected_satisfied(self._submitted_counts.items()) + + def can_accept(self, activity_history_id: str) -> bool: + expected_total = self._expected_counts.get(activity_history_id, 0) + if expected_total == 0: + return False + return self._submitted_counts.get(activity_history_id, 0) < expected_total + + def completion_state_after_add(self, activity_history_id: str) -> bool: + temp_counts = self._submitted_counts.copy() + temp_counts[activity_history_id] += 1 + return self._all_expected_satisfied(temp_counts.items()) + + def contains_activity(self, activity_history_id: str) -> bool: + if not self._has_flow_history: + return False + return activity_history_id in self._expected_ids + + @property + def expected_total(self) -> int: + return sum(self._expected_counts.values()) + + @property + def submitted_total(self) -> int: + return sum(self._submitted_counts.values()) + + def _all_expected_satisfied(self, submitted_items: Iterable[tuple[str, int]]) -> bool: + submitted_map = dict(submitted_items) + for activity_history_id, expected_count in self._expected_counts.items(): + if submitted_map.get(activity_history_id, 0) < expected_count: + return False + return True diff --git a/src/apps/answers/service.py b/src/apps/answers/service.py index 1a6815f66a9e..1d0e2736ff99 100644 --- a/src/apps/answers/service.py +++ b/src/apps/answers/service.py @@ -6,7 +6,7 @@ import os import time import uuid -from collections import Counter, defaultdict +from collections import defaultdict from json import JSONDecodeError from typing import Callable, List, Mapping, Optional @@ -89,6 +89,7 @@ WrongRespondentForAnswerGroup, ) from apps.answers.filters import AppletSubmitDateFilter, ReviewAppletItemFilter, SummaryActivityFilter +from apps.answers.flow_submission_progress import FlowSubmissionProgress from apps.answers.tasks import create_report from apps.applets.crud import AppletsCRUD from apps.applets.domain.applet_history import Version @@ -150,46 +151,38 @@ async def create_answer(self, activity_answer: AppletAnswerCreate, device_id: st async def _create_respondent_answer( self, activity_answer: AppletAnswerCreate, device_id: str | None ) -> AnswerSchema: - await self._validate_respondent_answer(activity_answer) - return await self._create_answer(activity_answer, device_id) + flow_progress = await self._validate_respondent_answer(activity_answer) + return await self._create_answer(activity_answer, device_id, flow_progress) async def _create_anonymous_answer( self, activity_answer: AppletAnswerCreate, device_id: str | None ) -> AnswerSchema: - await self._validate_anonymous_answer(activity_answer) - return await self._create_answer(activity_answer, device_id) + flow_progress = await self._validate_anonymous_answer(activity_answer) + return await self._create_answer(activity_answer, device_id, flow_progress) - async def _validate_respondent_answer(self, activity_answer: AppletAnswerCreate) -> None: - await self._validate_answer(activity_answer) + async def _validate_respondent_answer(self, activity_answer: AppletAnswerCreate) -> FlowSubmissionProgress | None: + flow_progress = await self._validate_answer(activity_answer) await self._validate_applet_for_user_response(activity_answer.applet_id) - async def _validate_anonymous_answer(self, activity_answer: AppletAnswerCreate) -> None: + return flow_progress + + async def _validate_anonymous_answer(self, activity_answer: AppletAnswerCreate) -> FlowSubmissionProgress | None: await self._validate_applet_for_anonymous_response(activity_answer.applet_id, activity_answer.version) - await self._validate_answer(activity_answer) + return await self._validate_answer(activity_answer) - async def _validate_answer(self, applet_answer: AppletAnswerCreate) -> None: # noqa: C901 + async def _validate_answer(self, applet_answer: AppletAnswerCreate) -> FlowSubmissionProgress | None: # noqa: C901 pk = self._generate_history_id(applet_answer.version) existed_answers = await AnswersCRUD(self.answer_session).get_by_submit_id(applet_answer.submit_id) activity_history_id = pk(applet_answer.activity_id) flow_history_id = pk(applet_answer.flow_id) if applet_answer.flow_id else None - - activity_indexes = set() # same activity is allowed multiple times in flow - latest_activity_index = None + flow_progress: FlowSubmissionProgress | None = None if flow_history_id: - flow_histories = await FlowsHistoryCRUD(self.session).load_full( - [pk(applet_answer.flow_id)], load_activities=False - ) - if not flow_histories: + flow_progress = FlowSubmissionProgress(self.session, self.answer_session) + await flow_progress.load(flow_history_id, applet_answer.submit_id) + if not flow_progress.has_flow_history: raise ValidationError("Flow not found") - flow_history = next(iter(flow_histories)) - - # check activity in the flow - for i, item in enumerate(flow_history.items): - if item.activity_id == activity_history_id: - activity_indexes.add(i) - latest_activity_index = len(flow_history.items) - 1 - if not activity_indexes: + if not flow_progress.contains_activity(activity_history_id): raise ValidationError("Activity not found in the flow") if existed_answers: @@ -208,60 +201,28 @@ async def _validate_answer(self, applet_answer: AppletAnswerCreate) -> None: # if flow_history_id != existed_answer.flow_history_id: raise ValidationError("Submit id duplicate error") - # check current answer is provided in right order in the flow, so prev activities already answered - prev_answers_count = len(existed_answers) - is_flow_completed = any(answer.is_flow_completed for answer in existed_answers) - - # Smart flow completion check - if is_flow_completed: - # Count expected and already persisted activity occurrences - flow_activity_counts = Counter(item.activity_id for item in flow_history.items) - submitted_counts = Counter(answer.activity_history_id for answer in existed_answers) - - # If all activities already persisted before this submission, the flow truly finished - if all(submitted_counts.get(act_id, 0) >= count for act_id, count in flow_activity_counts.items()): - raise ValidationError("Flow is already completed") + if flow_progress: + if flow_progress.is_complete_before_current(): + raise ValidationError("Flow is already completed") - current_expected_total = flow_activity_counts.get(activity_history_id, 0) - current_submitted = submitted_counts.get(activity_history_id, 0) - - # Reject duplicates that exceed expected occurrences for the activity - if current_submitted >= current_expected_total: - raise ValidationError("Flow is already completed") + if not flow_progress.can_accept(activity_history_id): + raise ValidationError("Activity submission exceeds expected occurrences for this flow") + if existed_answers: logger.info( - "Allowing late submission for flow %s, activity %s, submit_id %s", + "Allowing flow submission for flow %s, activity %s, submit_id %s", flow_history_id, activity_history_id, applet_answer.submit_id, ) - # Continue with existing order validation only if flow not marked complete - # When is_flow_completed=True, we allow out-of-order submissions - if not is_flow_completed and prev_answers_count not in activity_indexes: - assert latest_activity_index is not None - # allow latest activity for flow autocompletion FE logic - if not ( - prev_answers_count < latest_activity_index + 1 - and max(activity_indexes) == latest_activity_index - and applet_answer.is_flow_completed - ): - raise ValidationError("Wrong activity order in the flow") - - elif flow_history_id and 0 not in activity_indexes: - # check first flow answer - but allow if flow is marked as complete - # Check both existing answers and current answer for is_flow_completed - if not ( - (existed_answers and any(answer.is_flow_completed for answer in existed_answers)) - or applet_answer.is_flow_completed - ): - raise ValidationError("Wrong activity order in the flow") - activity_history = await ActivityHistoriesCRUD(self.session).get_by_id(activity_history_id) if not activity_history.applet_id.startswith(f"{applet_answer.applet_id}"): raise ActivityHistoryDoeNotExist() + return flow_progress + async def _validate_applet_for_anonymous_response(self, applet_id: uuid.UUID, version: str) -> None: await AppletHistoryService(self.session, applet_id, version).get() # Validate applet for anonymous answer @@ -346,12 +307,41 @@ async def _get_answer_relation( return relation.relation - async def _create_answer(self, applet_answer: AppletAnswerCreate, device_id: str | None) -> AnswerSchema: + async def _create_answer( + self, + applet_answer: AppletAnswerCreate, + device_id: str | None, + flow_progress: FlowSubmissionProgress | None, + ) -> AnswerSchema: assert self.user_id pk = self._generate_history_id(applet_answer.version) created_at = applet_answer.created_at or datetime.datetime.now(datetime.UTC).replace(tzinfo=None) subject_crud = SubjectsCrud(self.session) + activity_history_id = pk(applet_answer.activity_id) + flow_history_id = pk(applet_answer.flow_id) if applet_answer.flow_id else None + client_flow_completed_flag = bool(applet_answer.is_flow_completed) if applet_answer.flow_id else None + + is_flow_completed_backend = None + if applet_answer.flow_id: + if flow_progress: + is_flow_completed_backend = flow_progress.completion_state_after_add(activity_history_id) + if client_flow_completed_flag is not None and client_flow_completed_flag != is_flow_completed_backend: + logger.info( + "Flow completion mismatch for flow %s, activity %s, submit_id %s: client=%s backend=%s", + flow_history_id, + activity_history_id, + applet_answer.submit_id, + client_flow_completed_flag, + is_flow_completed_backend, + ) + else: + is_flow_completed_backend = client_flow_completed_flag + + migrated_data = None + if client_flow_completed_flag is not None: + migrated_data = {"client_flow_completed_flag": client_flow_completed_flag} + respondent_subject = await subject_crud.get_user_subject( user_id=self.user_id, applet_id=applet_answer.applet_id ) @@ -413,11 +403,11 @@ async def _create_answer(self, applet_answer: AppletAnswerCreate, device_id: str applet_id=applet_answer.applet_id, version=applet_answer.version, applet_history_id=pk(applet_answer.applet_id), - flow_history_id=pk(applet_answer.flow_id) if applet_answer.flow_id else None, - activity_history_id=pk(applet_answer.activity_id), + flow_history_id=flow_history_id, + activity_history_id=activity_history_id, respondent_id=self.user_id, client=applet_answer.client.dict(), - is_flow_completed=bool(applet_answer.is_flow_completed) if applet_answer.flow_id else None, + is_flow_completed=is_flow_completed_backend, target_subject_id=target_subject.id, source_subject_id=source_subject.id, input_subject_id=input_subject.id, @@ -425,6 +415,7 @@ async def _create_answer(self, applet_answer: AppletAnswerCreate, device_id: str consent_to_share=applet_answer.consent_to_share, event_history_id=applet_answer.event_history_id, device_id=device_id, + migrated_data=migrated_data, ) ) item_answer = applet_answer.answer diff --git a/src/apps/answers/tests/test_answers.py b/src/apps/answers/tests/test_answers.py index 8492b22ee245..83b6b5984691 100644 --- a/src/apps/answers/tests/test_answers.py +++ b/src/apps/answers/tests/test_answers.py @@ -503,12 +503,13 @@ async def tom_answer_activity_flow_incomplete( session: AsyncSession, tom: User, applet_with_flow: AppletFull ) -> AnswerSchema: answer_service = AnswerService(session, tom.id) + # Use flow[1] which has 3 activities and only submit the first one to create truly incomplete flow return await answer_service.create_answer( AppletAnswerCreate( applet_id=applet_with_flow.id, version=applet_with_flow.version, submit_id=uuid.uuid4(), - flow_id=applet_with_flow.activity_flows[0].id, + flow_id=applet_with_flow.activity_flows[1].id, is_flow_completed=False, activity_id=applet_with_flow.activities[0].id, answer=ItemAnswerCreate( @@ -551,7 +552,7 @@ def applet_with_flow_answer_create(applet_with_flow: AppletFull) -> list[AppletA **answer_data, consent_to_share=False, ), - # flow#2 submission#1 + # flow#2 submission#1 (flow[1] has 3 activities) AppletAnswerCreate( submit_id=submit_id, flow_id=applet_with_flow.activity_flows[1].id, @@ -568,12 +569,21 @@ def applet_with_flow_answer_create(applet_with_flow: AppletFull) -> list[AppletA AppletAnswerCreate( submit_id=submit_id, flow_id=applet_with_flow.activity_flows[1].id, - is_flow_completed=True, + is_flow_completed=False, activity_id=applet_with_flow.activities[1].id, answer=ItemAnswerCreate(item_ids=[applet_with_flow.activities[1].items[0].id], **answer_item_data), **answer_data, consent_to_share=False, ), + AppletAnswerCreate( + submit_id=submit_id, + flow_id=applet_with_flow.activity_flows[1].id, + is_flow_completed=True, + activity_id=applet_with_flow.activities[2].id, + answer=ItemAnswerCreate(item_ids=[applet_with_flow.activities[2].items[0].id], **answer_item_data), + **answer_data, + consent_to_share=False, + ), # flow#1 submission#2 AppletAnswerCreate( submit_id=uuid.uuid4(), @@ -1255,18 +1265,24 @@ async def test_create_flow_answer__wrong_order( data.flow_id = applet_with_flow.activity_flows[1].id data.activity_id = applet_with_flow.activity_flows[1].items[1].activity_id + # With our out-of-order fix, submitting the second activity first is now allowed response = await client.post(self.answer_url, data=data) - assert response.status_code == http.HTTPStatus.BAD_REQUEST + assert response.status_code == http.HTTPStatus.CREATED - # different submit_id for second activity + # Submit first activity with same submit_id data.activity_id = applet_with_flow.activity_flows[1].items[0].activity_id - response = await client.post(self.answer_url, data=data) assert response.status_code == http.HTTPStatus.CREATED - data.submit_id = uuid.uuid4() - data.activity_id = applet_with_flow.activity_flows[1].items[1].activity_id + # Submit third activity to complete the flow (flow[1] has 3 activities) + data.activity_id = applet_with_flow.activity_flows[1].items[2].activity_id + data.is_flow_completed = True + response = await client.post(self.answer_url, data=data) + assert response.status_code == http.HTTPStatus.CREATED + # Now trying to submit any activity with the SAME submit_id should fail (flow already complete) + # Keep the same submit_id to test duplicate submission + data.activity_id = applet_with_flow.activity_flows[1].items[0].activity_id response = await client.post(self.answer_url, data=data) assert response.status_code == http.HTTPStatus.BAD_REQUEST @@ -2379,8 +2395,9 @@ async def test_get_flow_identifiers_incomplete_submission( tom_subject = await SubjectsService(session, tom.id).get_by_user_and_applet(tom.id, applet.id) assert tom_subject + # Use flow[1] which is the incomplete flow from the fixture identifier_url = self.flow_identifiers_url.format( - applet_id=applet.id, flow_id=applet_with_flow.activity_flows[0].id + applet_id=applet.id, flow_id=applet_with_flow.activity_flows[1].id ) response = await client.get(identifier_url, dict(targetSubjectId=tom_subject.id)) @@ -3035,7 +3052,7 @@ async def test_flow_submission_incomplete( client.login(tom) url = self.flow_submission_url.format( applet_id=applet_with_flow.id, - flow_id=applet_with_flow.activity_flows[0].id, + flow_id=applet_with_flow.activity_flows[1].id, # Use flow[1] which matches the fixture submit_id=tom_answer_activity_flow_incomplete.submit_id, ) response = await client.get(url) @@ -3045,7 +3062,8 @@ async def test_flow_submission_incomplete( data = data["result"] assert set(data.keys()) == {"flow", "submission", "summary"} assert data["submission"]["isCompleted"] is False - assert len(data["submission"]["answers"]) == len(applet_with_flow.activity_flows[0].items) + # Only 1 answer submitted out of 3 activities in flow[1] + assert len(data["submission"]["answers"]) == 1 async def test_flow_submission_no_flow( self, client, tom: User, applet_with_flow: AppletFull, tom_answer_activity_no_flow @@ -3171,7 +3189,7 @@ async def test_get_flow_submissions_incomplete( client.login(tom) url = self.flow_submissions_url.format( applet_id=applet_with_flow.id, - flow_id=applet_with_flow.activity_flows[0].id, + flow_id=applet_with_flow.activity_flows[1].id, # Use flow[1] which matches the fixture ) tom_subject = await SubjectsService(session, tom.id).get_by_user_and_applet(tom.id, applet_with_flow.id) diff --git a/src/apps/answers/tests/test_flow_out_of_order.py b/src/apps/answers/tests/test_flow_out_of_order.py index 80dcb94d394e..cb5e1cfd8130 100644 --- a/src/apps/answers/tests/test_flow_out_of_order.py +++ b/src/apps/answers/tests/test_flow_out_of_order.py @@ -5,6 +5,7 @@ from sqlalchemy.ext.asyncio import AsyncSession +from apps.answers.crud.answers import AnswersCRUD from apps.answers.domain import AppletAnswerCreate, ClientMeta, ItemAnswerCreate from apps.answers.service import AnswerService from apps.applets.domain.applet_full import AppletFull @@ -61,6 +62,37 @@ async def test_out_of_order_submission_accepts_all_activities( response = await client.post(self.answer_url, data=data_b) assert response.status_code == http.HTTPStatus.CREATED, f"Failed to create activity B: {response.json()}" + async def test_backend_marks_only_last_activity_complete_for_legacy_clients( + self, + client: TestClient, + tom: User, + answer_create: AppletAnswerCreate, + applet_with_flow: AppletFull, + session: AsyncSession, + ): + """Legacy clients set completion on every activity; backend should mark only the last one.""" + client.login(tom) + flow = applet_with_flow.activity_flows[1] + submit_id = uuid.uuid4() + + for item in flow.items: + data = answer_create.copy(deep=True) + data.submit_id = submit_id + data.applet_id = applet_with_flow.id + data.flow_id = flow.id + data.activity_id = item.activity_id + data.is_flow_completed = True # Legacy behaviour + + response = await client.post(self.answer_url, data=data) + assert response.status_code == http.HTTPStatus.CREATED, response.json() + + stored_answers = await AnswersCRUD(session).get_by_submit_id(submit_id) + assert stored_answers is not None + assert len(stored_answers) == len(flow.items) + backend_flags = [answer.is_flow_completed for answer in stored_answers] + assert all(flag is False for flag in backend_flags[:-1]) + assert backend_flags[-1] is True + async def test_duplicate_activity_rejected_after_flow_complete( self, client: TestClient, @@ -102,18 +134,18 @@ async def test_flow_with_duplicate_activities_handles_counts_correctly( client: TestClient, tom: User, answer_create: AppletAnswerCreate, - applet_with_flow: AppletFull, + applet_with_flow_duplicated_activities: AppletFull, ): """Test flow where the same activity appears multiple times""" client.login(tom) - flow = applet_with_flow.activity_flows[0] # Flow with duplicate activities + flow = applet_with_flow_duplicated_activities.activity_flows[0] # Flow with duplicate activities submit_id = uuid.uuid4() # Assume flow has activity A twice (at positions 0 and 2) # Submit first occurrence data1 = answer_create.copy(deep=True) data1.submit_id = submit_id - data1.applet_id = applet_with_flow.id + data1.applet_id = applet_with_flow_duplicated_activities.id data1.flow_id = flow.id data1.activity_id = flow.items[0].activity_id data1.is_flow_completed = False @@ -124,15 +156,25 @@ async def test_flow_with_duplicate_activities_handles_counts_correctly( # Submit second occurrence should still be accepted data2 = answer_create.copy(deep=True) data2.submit_id = submit_id - data2.applet_id = applet_with_flow.id + data2.applet_id = applet_with_flow_duplicated_activities.id data2.flow_id = flow.id data2.activity_id = flow.items[0].activity_id # Same activity ID data2.is_flow_completed = True response = await client.post(self.answer_url, data=data2) - # This test assumes the flow actually has duplicates - adjust based on test fixtures - # For now, we'll check that it's handled without error - assert response.status_code in [http.HTTPStatus.CREATED, http.HTTPStatus.BAD_REQUEST] + assert response.status_code == http.HTTPStatus.CREATED + + # Third occurrence should be rejected as it exceeds expected count + data3 = answer_create.copy(deep=True) + data3.submit_id = submit_id + data3.applet_id = applet_with_flow_duplicated_activities.id + data3.flow_id = flow.id + data3.activity_id = flow.items[0].activity_id + data3.is_flow_completed = False + + response = await client.post(self.answer_url, data=data3) + assert response.status_code == http.HTTPStatus.BAD_REQUEST + assert "Flow is already completed" in response.json()["result"][0]["message"] @patch("apps.answers.service.logger") async def test_late_submission_logs_correctly( @@ -190,7 +232,7 @@ async def test_late_submission_logs_correctly( # Verify logging was called mock_logger.info.assert_called() log_call = mock_logger.info.call_args[0][0] - assert "Allowing late submission" in log_call + assert "Allowing flow submission" in log_call async def test_normal_flow_behavior_unchanged( self, From 93bc172a7d5bf9891a360c93552441d15e5ab608 Mon Sep 17 00:00:00 2001 From: Mohamed Elmi Date: Wed, 8 Oct 2025 02:49:50 -0500 Subject: [PATCH 2/4] fix: add back create_at and convert Unix timestamp to datetime in answer existence check --- src/apps/answers/api.py | 4 ++-- src/apps/answers/crud/answers.py | 10 ++++++++-- src/apps/answers/domain/answers.py | 6 ++---- src/apps/answers/service.py | 4 ++-- 4 files changed, 14 insertions(+), 10 deletions(-) diff --git a/src/apps/answers/api.py b/src/apps/answers/api.py index 10d6ddd25901..f12040644665 100644 --- a/src/apps/answers/api.py +++ b/src/apps/answers/api.py @@ -873,12 +873,12 @@ async def answers_existence_check( await AppletService(session, user.id).exist_by_id(schema.applet_id) await CheckAccessService(session, user.id).check_answer_check_access(schema.applet_id) is_exist = await AnswerService(session, user.id, answer_session).is_answers_uploaded( - schema.applet_id, schema.activity_id, schema.submit_id + schema.applet_id, schema.activity_id, schema.submit_id, schema.created_at ) logger.info( f"check-existence: applet_id={schema.applet_id}, activity_id={schema.activity_id}, user_id={user.id}, " - f"submit_id={schema.submit_id}, exists={is_exist}, ip={client_ip}" + f"submit_id={schema.submit_id}, created_at={schema.created_at}, exists={is_exist}, ip={client_ip}" ) return Response[AnswerExistenceResponse](result=AnswerExistenceResponse(exists=is_exist)) diff --git a/src/apps/answers/crud/answers.py b/src/apps/answers/crud/answers.py index 307176f699b4..f04ce35d1519 100644 --- a/src/apps/answers/crud/answers.py +++ b/src/apps/answers/crud/answers.py @@ -518,9 +518,9 @@ async def get_by_applet_activity_submit_or_user_id( activity_id: str, user_id: uuid.UUID | None = None, submit_id: uuid.UUID | None = None, + created_at: int | None = None, ) -> list[AnswerSchema]: - # We're not using created_at for filtering as it causes issues with mobile submissions - # The combination of applet_id, activity_id, and either user_id or submit_id should be sufficient + # created_at is used to distinguish between duplicate activities in flows query: Query = select(AnswerSchema) query = query.where(AnswerSchema.applet_id == applet_id) query = query.filter(AnswerSchema.activity_history_id.startswith(activity_id)) @@ -528,6 +528,12 @@ async def get_by_applet_activity_submit_or_user_id( query = query.where(AnswerSchema.submit_id == submit_id) if user_id: query = query.where(AnswerSchema.respondent_id == user_id) + if created_at is not None: + # Convert Unix timestamp (milliseconds) to datetime for comparison + created_at_datetime = datetime.datetime.fromtimestamp( + created_at / 1000.0, tz=datetime.timezone.utc + ).replace(tzinfo=None) + query = query.where(AnswerSchema.created_at == created_at_datetime) db_result = await self._execute(query) return db_result.scalars().all() diff --git a/src/apps/answers/domain/answers.py b/src/apps/answers/domain/answers.py index 63b7fa44e7f7..5cd112c06ed8 100644 --- a/src/apps/answers/domain/answers.py +++ b/src/apps/answers/domain/answers.py @@ -673,10 +673,8 @@ class AppletCompletedEntities(InternalModel): class AnswersCheck(PublicModel): applet_id: uuid.UUID - # TODO: created_at can be safely removed after - # the corresponding mobile PR is merged - # https://mindlogger.atlassian.net/browse/M2-9693 - created_at: int | None = None + # Used to distinguish between duplicate activities in flows + created_at: int activity_id: str submit_id: uuid.UUID | None = None diff --git a/src/apps/answers/service.py b/src/apps/answers/service.py index 1d0e2736ff99..d63813e31767 100644 --- a/src/apps/answers/service.py +++ b/src/apps/answers/service.py @@ -1715,11 +1715,11 @@ async def get_completed_answers_data_list( return result async def is_answers_uploaded( - self, applet_id: uuid.UUID, activity_id: str, submit_id: uuid.UUID | None = None + self, applet_id: uuid.UUID, activity_id: str, submit_id: uuid.UUID | None = None, created_at: int | None = None ) -> bool: # check by submit id if provided otherwise by user_id answers = await AnswersCRUD(self.answer_session).get_by_applet_activity_submit_or_user_id( - applet_id, activity_id, self.user_id if not submit_id else None, submit_id + applet_id, activity_id, self.user_id if not submit_id else None, submit_id, created_at ) if not answers: return False From f08d77565168643cef17232d3274db521140a068 Mon Sep 17 00:00:00 2001 From: Mohamed Elmi Date: Mon, 13 Oct 2025 14:26:18 -0500 Subject: [PATCH 3/4] M2-9862 make the create_at optional --- src/apps/answers/domain/answers.py | 2 +- src/apps/answers/service.py | 18 ++- src/apps/answers/tests/conftest.py | 2 +- .../answers/tests/test_flow_out_of_order.py | 150 +++++++++++++++++- 4 files changed, 165 insertions(+), 7 deletions(-) diff --git a/src/apps/answers/domain/answers.py b/src/apps/answers/domain/answers.py index 5cd112c06ed8..4671e4f8af88 100644 --- a/src/apps/answers/domain/answers.py +++ b/src/apps/answers/domain/answers.py @@ -674,7 +674,7 @@ class AppletCompletedEntities(InternalModel): class AnswersCheck(PublicModel): applet_id: uuid.UUID # Used to distinguish between duplicate activities in flows - created_at: int + created_at: int | None = None activity_id: str submit_id: uuid.UUID | None = None diff --git a/src/apps/answers/service.py b/src/apps/answers/service.py index d63813e31767..9e26f7d14ae4 100644 --- a/src/apps/answers/service.py +++ b/src/apps/answers/service.py @@ -172,12 +172,28 @@ async def _validate_anonymous_answer(self, activity_answer: AppletAnswerCreate) async def _validate_answer(self, applet_answer: AppletAnswerCreate) -> FlowSubmissionProgress | None: # noqa: C901 pk = self._generate_history_id(applet_answer.version) + + # Timestamp-based duplicate detection: when created_at provided, check for exact duplicate + # Each unique timestamp represents a distinct submission, bypassing occurrence limits + if applet_answer.created_at is not None: + created_at_ms = int(applet_answer.created_at.timestamp() * 1000) + existing_with_timestamp = await AnswersCRUD(self.answer_session).get_by_applet_activity_submit_or_user_id( + applet_answer.applet_id, + str(applet_answer.activity_id), + None, + applet_answer.submit_id, + created_at_ms, + ) + if existing_with_timestamp: + raise ValidationError("Duplicate answer with same timestamp already exists") + existed_answers = await AnswersCRUD(self.answer_session).get_by_submit_id(applet_answer.submit_id) activity_history_id = pk(applet_answer.activity_id) flow_history_id = pk(applet_answer.flow_id) if applet_answer.flow_id else None flow_progress: FlowSubmissionProgress | None = None - if flow_history_id: + # Only use occurrence-based flow validation when created_at is NOT provided + if flow_history_id and applet_answer.created_at is None: flow_progress = FlowSubmissionProgress(self.session, self.answer_session) await flow_progress.load(flow_history_id, applet_answer.submit_id) if not flow_progress.has_flow_history: diff --git a/src/apps/answers/tests/conftest.py b/src/apps/answers/tests/conftest.py index b3a8e7c058a0..c18802072ec8 100644 --- a/src/apps/answers/tests/conftest.py +++ b/src/apps/answers/tests/conftest.py @@ -235,7 +235,7 @@ def answer_create( submit_id=uuid.uuid4(), activity_id=applet.activities[0].id, answer=answer_item_create, - created_at=datetime.datetime.now(datetime.UTC).replace(microsecond=0), + created_at=None, # None by default - uses occurrence-based validation client=client_meta, consent_to_share=False, ) diff --git a/src/apps/answers/tests/test_flow_out_of_order.py b/src/apps/answers/tests/test_flow_out_of_order.py index cb5e1cfd8130..a71739e34bec 100644 --- a/src/apps/answers/tests/test_flow_out_of_order.py +++ b/src/apps/answers/tests/test_flow_out_of_order.py @@ -184,12 +184,12 @@ async def test_late_submission_logs_correctly( tom: User, applet_with_flow: AppletFull, ): - """Test that late submissions are logged for monitoring""" + """Test that late submissions are logged for monitoring (occurrence-based validation)""" service = AnswerService(session, tom.id) flow = applet_with_flow.activity_flows[1] submit_id = uuid.uuid4() - # Create last activity with is_flow_completed=True + # Create last activity with is_flow_completed=True (without created_at for occurrence-based validation) answer_c = AppletAnswerCreate( applet_id=applet_with_flow.id, version=applet_with_flow.version, @@ -203,7 +203,7 @@ async def test_late_submission_logs_correctly( end_time=datetime.datetime.now(datetime.UTC), user_public_key="test_key", ), - created_at=datetime.datetime.now(datetime.UTC), + created_at=None, # Use occurrence-based validation client=ClientMeta(app_id="test", app_version="1.0.0"), ) @@ -223,7 +223,7 @@ async def test_late_submission_logs_correctly( end_time=datetime.datetime.now(datetime.UTC), user_public_key="test_key", ), - created_at=datetime.datetime.now(datetime.UTC), + created_at=None, # Use occurrence-based validation client=ClientMeta(app_id="test", app_version="1.0.0"), ) @@ -324,3 +324,145 @@ async def test_partial_out_of_order_with_missing_activities( response = await client.post(self.answer_url, data=duplicate_data) assert response.status_code == http.HTTPStatus.BAD_REQUEST assert "Flow is already completed" in response.json()["result"][0]["message"] + + async def test_timestamp_allows_duplicate_activities_with_different_timestamps( + self, + client: TestClient, + tom: User, + answer_create: AppletAnswerCreate, + applet_with_flow: AppletFull, + ): + """When created_at is provided, different timestamps allow duplicate activities""" + client.login(tom) + flow = applet_with_flow.activity_flows[1] + submit_id = uuid.uuid4() + + # Submit same activity with first timestamp + data1 = answer_create.copy(deep=True) + data1.submit_id = submit_id + data1.applet_id = applet_with_flow.id + data1.flow_id = flow.id + data1.activity_id = flow.items[0].activity_id + data1.created_at = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) + data1.is_flow_completed = False + + response = await client.post(self.answer_url, data=data1) + assert response.status_code == http.HTTPStatus.CREATED, f"First submission failed: {response.json()}" + + # Submit same activity with different timestamp (should be accepted) + data2 = answer_create.copy(deep=True) + data2.submit_id = submit_id + data2.applet_id = applet_with_flow.id + data2.flow_id = flow.id + data2.activity_id = flow.items[0].activity_id + data2.created_at = datetime.datetime(2024, 1, 1, 12, 30, 0, tzinfo=datetime.UTC) + data2.is_flow_completed = False + + response = await client.post(self.answer_url, data=data2) + assert response.status_code == http.HTTPStatus.CREATED, f"Second submission with different timestamp failed: {response.json()}" + + async def test_timestamp_rejects_duplicate_with_same_timestamp( + self, + client: TestClient, + tom: User, + answer_create: AppletAnswerCreate, + applet_with_flow: AppletFull, + ): + """When created_at is provided, same timestamp rejects duplicate""" + client.login(tom) + flow = applet_with_flow.activity_flows[1] + submit_id = uuid.uuid4() + timestamp = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) + + # Submit activity with timestamp + data1 = answer_create.copy(deep=True) + data1.submit_id = submit_id + data1.applet_id = applet_with_flow.id + data1.flow_id = flow.id + data1.activity_id = flow.items[0].activity_id + data1.created_at = timestamp + data1.is_flow_completed = False + + response = await client.post(self.answer_url, data=data1) + assert response.status_code == http.HTTPStatus.CREATED + + # Try to submit same activity with same timestamp (should be rejected) + data2 = answer_create.copy(deep=True) + data2.submit_id = submit_id + data2.applet_id = applet_with_flow.id + data2.flow_id = flow.id + data2.activity_id = flow.items[0].activity_id + data2.created_at = timestamp + data2.is_flow_completed = False + + response = await client.post(self.answer_url, data=data2) + assert response.status_code == http.HTTPStatus.BAD_REQUEST + assert "Duplicate answer with same timestamp" in response.json()["result"][0]["message"] + + async def test_timestamp_bypasses_occurrence_limits( + self, + client: TestClient, + tom: User, + answer_create: AppletAnswerCreate, + applet_with_flow: AppletFull, + ): + """When created_at is provided, occurrence counting is bypassed""" + client.login(tom) + flow = applet_with_flow.activity_flows[1] + submit_id = uuid.uuid4() + + # Submit same activity multiple times with different timestamps + # This would normally be rejected by occurrence counting + for i in range(5): + data = answer_create.copy(deep=True) + data.submit_id = submit_id + data.applet_id = applet_with_flow.id + data.flow_id = flow.id + data.activity_id = flow.items[0].activity_id + data.created_at = datetime.datetime(2024, 1, 1, 12, i, 0, tzinfo=datetime.UTC) + data.is_flow_completed = False + + response = await client.post(self.answer_url, data=data) + assert response.status_code == http.HTTPStatus.CREATED, f"Submission {i} failed: {response.json()}" + + async def test_check_existence_consistent_with_upload_validation( + self, + client: TestClient, + tom: User, + answer_create: AppletAnswerCreate, + applet_with_flow: AppletFull, + ): + """Check-existence endpoint should match upload validation behavior""" + client.login(tom) + flow = applet_with_flow.activity_flows[1] + submit_id = uuid.uuid4() + timestamp = datetime.datetime(2024, 1, 1, 12, 0, 0, tzinfo=datetime.UTC) + + # Submit activity with timestamp + data = answer_create.copy(deep=True) + data.submit_id = submit_id + data.applet_id = applet_with_flow.id + data.flow_id = flow.id + data.activity_id = flow.items[0].activity_id + data.created_at = timestamp + data.is_flow_completed = False + + response = await client.post(self.answer_url, data=data) + assert response.status_code == http.HTTPStatus.CREATED + + # Check existence with same timestamp (should return true) + check_data = { + "applet_id": str(applet_with_flow.id), + "activity_id": str(flow.items[0].activity_id), + "submit_id": str(submit_id), + "created_at": int(timestamp.timestamp() * 1000), + } + response = await client.post("/answers/check-existence", data=check_data) + assert response.status_code == http.HTTPStatus.OK + assert response.json()["result"]["exists"] is True + + # Check existence with different timestamp (should return false) + check_data["created_at"] = int(datetime.datetime(2024, 1, 1, 13, 0, 0, tzinfo=datetime.UTC).timestamp() * 1000) + response = await client.post("/answers/check-existence", data=check_data) + assert response.status_code == http.HTTPStatus.OK + assert response.json()["result"]["exists"] is False From f8d62701ec1394229b79308951baa82b501d5122 Mon Sep 17 00:00:00 2001 From: Mohamed Elmi Date: Fri, 17 Oct 2025 10:14:29 -0500 Subject: [PATCH 4/4] fix(M2-9862): fix linting using ruff --- src/apps/answers/tests/test_flow_out_of_order.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/apps/answers/tests/test_flow_out_of_order.py b/src/apps/answers/tests/test_flow_out_of_order.py index a71739e34bec..b0b0db4da72d 100644 --- a/src/apps/answers/tests/test_flow_out_of_order.py +++ b/src/apps/answers/tests/test_flow_out_of_order.py @@ -359,7 +359,9 @@ async def test_timestamp_allows_duplicate_activities_with_different_timestamps( data2.is_flow_completed = False response = await client.post(self.answer_url, data=data2) - assert response.status_code == http.HTTPStatus.CREATED, f"Second submission with different timestamp failed: {response.json()}" + assert response.status_code == http.HTTPStatus.CREATED, ( + f"Second submission with different timestamp failed: {response.json()}" + ) async def test_timestamp_rejects_duplicate_with_same_timestamp( self,