diff --git a/openapi/specs/b2b_learner_records.yaml b/openapi/specs/b2b_learner_records.yaml index d77ae68..cf3a410 100644 --- a/openapi/specs/b2b_learner_records.yaml +++ b/openapi/specs/b2b_learner_records.yaml @@ -53,6 +53,19 @@ paths: description: Return only records changed at or after this instant. title: Updated Since description: Return only records changed at or after this instant. + - name: updated_before + in: query + required: false + schema: + anyOf: + - type: string + format: date-time + - type: 'null' + description: Return only records changed strictly before this instant. Paired + with updated_since, bounds a window for a partitioned backfill or replay. + title: Updated Before + description: Return only records changed strictly before this instant. Paired + with updated_since, bounds a window for a partitioned backfill or replay. - name: include_inactive in: query required: false @@ -156,6 +169,19 @@ paths: description: Return only records changed at or after this instant. title: Updated Since description: Return only records changed at or after this instant. + - name: updated_before + in: query + required: false + schema: + anyOf: + - type: string + format: date-time + - type: 'null' + description: Return only records changed strictly before this instant. Paired + with updated_since, bounds a window for a partitioned backfill or replay. + title: Updated Before + description: Return only records changed strictly before this instant. Paired + with updated_since, bounds a window for a partitioned backfill or replay. - name: include_inactive in: query required: false @@ -220,6 +246,50 @@ paths: - type: integer - type: 'null' title: Contract Id + - name: contract_is_active + in: query + required: false + schema: + anyOf: + - type: boolean + - type: 'null' + title: Contract Is Active + - name: courserun_id + in: query + required: false + schema: + anyOf: + - type: string + - type: 'null' + title: Courserun Id + - name: courserun_starts_after + in: query + required: false + schema: + anyOf: + - type: string + format: date-time + - type: 'null' + description: Return only course runs starting at or after this instant. + Self-paced or otherwise unscheduled runs have no start date and match + neither this nor courserun_starts_before; a partitioned sync needs one + pass with neither set to pick those up. + title: Courserun Starts After + description: Return only course runs starting at or after this instant. Self-paced + or otherwise unscheduled runs have no start date and match neither this + nor courserun_starts_before; a partitioned sync needs one pass with neither + set to pick those up. + - name: courserun_starts_before + in: query + required: false + schema: + anyOf: + - type: string + format: date-time + - type: 'null' + description: Return only course runs starting strictly before this instant. + title: Courserun Starts Before + description: Return only course runs starting strictly before this instant. - name: limit in: query required: false diff --git a/src/ol_analytics_api/tenants/b2b_learner_records/queries.py b/src/ol_analytics_api/tenants/b2b_learner_records/queries.py index 2a7005f..45cb19f 100644 --- a/src/ol_analytics_api/tenants/b2b_learner_records/queries.py +++ b/src/ol_analytics_api/tenants/b2b_learner_records/queries.py @@ -114,10 +114,14 @@ def _outcomes_shared() -> str: class RecordFilters: organization_id: uuid.UUID contract_id: int | None = None + contract_is_active: bool | None = None courserun_id: str | None = None + courserun_starts_after: datetime.datetime | None = None + courserun_starts_before: datetime.datetime | None = None learner_ids: tuple[uuid.UUID, ...] = () completion_statuses: tuple[str, ...] = () updated_since: datetime.datetime | None = None + updated_before: datetime.datetime | None = None include_inactive: bool = False @@ -137,17 +141,32 @@ def _placeholders(count: int) -> str: return ", ".join(["%s"] * count) -def _cursor_value(value: datetime.datetime) -> str: - """Render ``updated_since`` for comparison against ``record_updated_on``. - - The MVs store that cursor as a zone-less UTC ISO-8601 string, so the - comparison is lexicographic. Truncating to whole seconds makes the bound a - prefix of any stored value in the same second, whatever its fractional - precision, so a record at the boundary is re-sent rather than skipped. +def _cursor_value(value: datetime.datetime, *, round_up: bool = False) -> str: + """Render a datetime for comparison against an MV timestamp column. + + The MVs store every MITx Online timestamp as a zone-less UTC ISO-8601 + string, so the comparison is lexicographic and only whole-second precision + survives. For an inclusive lower bound (the default, ``round_up=False``), + truncating down makes the bound a prefix of any stored value in the same + second, so a record at the boundary is matched by ``>=`` rather than + skipped — the record may be re-sent by the next window too, which a + sync client tolerates, but it is never dropped. + + ``round_up=True`` is the same trade for an exclusive upper bound: rounding + a fractional instant *down* would turn ``< bound`` into ``< floor(bound)``, + which excludes every record in that second, including ones genuinely + before the requested instant — a silent drop, not a re-send. Rounding up + to the next whole second instead means ``<`` matches everything through + that second, so the same record may be re-sent by the window that starts + there (its ``since`` floors to the same second), never lost between them. """ if value.tzinfo is not None: value = value.astimezone(datetime.UTC).replace(tzinfo=None) - return value.replace(microsecond=0).isoformat() + if round_up and value.microsecond: + value = value.replace(microsecond=0) + datetime.timedelta(seconds=1) + else: + value = value.replace(microsecond=0) + return value.isoformat() def _assemble( # noqa: PLR0913, PLR0917 @@ -181,15 +200,38 @@ def _assemble( # noqa: PLR0913, PLR0917 return RecordQuery(page, count, (*record_params, *predicate_params), sources) +def _window_predicates( + column: str, since: datetime.datetime | None, before: datetime.datetime | None +) -> tuple[list[str], list[Any]]: + """``[since, before)``: paired with a later window whose ``since`` is this + ``before``, a backfill can hand each window to one worker and know no row + is dropped between them. Both bounds round toward re-sending a row rather + than skipping it, so a row at a sub-second boundary can land in both + windows, never in neither. ``column`` is always one of this module's own + fixed identifiers, never caller input. + """ + predicates: list[str] = [] + params: list[Any] = [] + if since is not None: + predicates.append(f"{column} >= %s") + params.append(_cursor_value(since)) + if before is not None: + predicates.append(f"{column} < %s") + params.append(_cursor_value(before, round_up=True)) + return predicates, params + + def _shared_predicates(filters: RecordFilters) -> tuple[list[str], list[Any]]: predicates: list[str] = [] params: list[Any] = [] if filters.learner_ids: predicates.append(f"learner_id IN ({_placeholders(len(filters.learner_ids))})") params.extend(str(learner_id) for learner_id in filters.learner_ids) - if filters.updated_since is not None: - predicates.append("record_updated_on >= %s") - params.append(_cursor_value(filters.updated_since)) + window_predicates, window_params = _window_predicates( + "record_updated_on", filters.updated_since, filters.updated_before + ) + predicates.extend(window_predicates) + params.extend(window_params) return predicates, params @@ -384,6 +426,14 @@ def courses(schema: str, filters: RecordFilters) -> RecordQuery: Built without ``_assemble``: these rows carry no personal data, so there is no consent projection, and ``outcomes_withheld_count`` is always 0. + + This MV carries no ``record_updated_on``, so ``courserun_starts_after``/ + ``courserun_starts_before`` partition by when a run starts, not by when its + row last changed. ``courserun_start_on`` is nullable (unscheduled runs), and + SQL comparisons against NULL are never true, so a self-paced run with no + start date matches neither bound and falls outside every partition. A + partitioned backfill that wants full coverage still needs one unfiltered + pass (or a pass with neither bound set) to pick those up. """ table = f"{validate_sql_identifier(schema)}.{CONTRACT_COURSERUN_MV}" scope = ["sso_organization_id = %s"] @@ -391,6 +441,17 @@ def courses(schema: str, filters: RecordFilters) -> RecordQuery: if filters.contract_id is not None: scope.append("contract_id = %s") params.append(filters.contract_id) + if filters.contract_is_active is not None: + scope.append("b2b_contract_is_active = %s") + params.append(filters.contract_is_active) + if filters.courserun_id is not None: + scope.append("courserun_readable_id = %s") + params.append(filters.courserun_id) + window_scope, window_params = _window_predicates( + "courserun_start_on", filters.courserun_starts_after, filters.courserun_starts_before + ) + scope.extend(window_scope) + params.extend(window_params) where = " AND ".join(scope) page = ( "SELECT sso_organization_id AS organization_id, organization_name, contract_id," # noqa: S608 diff --git a/src/ol_analytics_api/tenants/b2b_learner_records/routers/organizations.py b/src/ol_analytics_api/tenants/b2b_learner_records/routers/organizations.py index 5d06ffb..fc46102 100644 --- a/src/ol_analytics_api/tenants/b2b_learner_records/routers/organizations.py +++ b/src/ol_analytics_api/tenants/b2b_learner_records/routers/organizations.py @@ -64,6 +64,15 @@ def page( datetime.datetime | None, Query(description="Return only records changed at or after this instant."), ] +UpdatedBefore = Annotated[ + datetime.datetime | None, + Query( + description=( + "Return only records changed strictly before this instant. Paired with " + "updated_since, bounds a window for a partitioned backfill or replay." + ) + ), +] IncludeInactive = Annotated[ bool, Query(description="Include deactivated enrollments (unenrolled, refunded, transferred).") ] @@ -109,6 +118,7 @@ async def list_learners( # noqa: PLR0913 contract_id: int | None = None, learner_id: LearnerIds, updated_since: UpdatedSince = None, + updated_before: UpdatedBefore = None, include_inactive: IncludeInactive = False, ) -> LearnerRecordsResponse[BaseModel]: filters = queries.RecordFilters( @@ -116,6 +126,7 @@ async def list_learners( # noqa: PLR0913 contract_id=contract_id, learner_ids=tuple(learner_id or ()), updated_since=updated_since, + updated_before=updated_before, include_inactive=include_inactive, ) return await _respond( @@ -139,6 +150,7 @@ async def list_enrollments( # noqa: PLR0913 learner_id: LearnerIds, completion_status: Annotated[list[CompletionStatusFilter], Query(default_factory=list)], updated_since: UpdatedSince = None, + updated_before: UpdatedBefore = None, include_inactive: IncludeInactive = False, ) -> LearnerRecordsResponse[BaseModel]: filters = queries.RecordFilters( @@ -148,6 +160,7 @@ async def list_enrollments( # noqa: PLR0913 learner_ids=tuple(learner_id or ()), completion_statuses=tuple(status.value for status in completion_status or ()), updated_since=updated_since, + updated_before=updated_before, include_inactive=include_inactive, ) return await _respond( @@ -162,13 +175,37 @@ async def list_enrollments( # noqa: PLR0913 response_model=LearnerRecordsResponse[CourseRun], summary="Contracts and course runs covered by the organization's licence", ) -async def list_courses( +async def list_courses( # noqa: PLR0913 *, organization_id: uuid.UUID, page: PageParams, contract_id: int | None = None, + contract_is_active: bool | None = None, + courserun_id: str | None = None, + courserun_starts_after: Annotated[ + datetime.datetime | None, + Query( + description=( + "Return only course runs starting at or after this instant. Self-paced or " + "otherwise unscheduled runs have no start date and match neither this nor " + "courserun_starts_before; a partitioned sync needs one pass with neither " + "set to pick those up." + ) + ), + ] = None, + courserun_starts_before: Annotated[ + datetime.datetime | None, + Query(description="Return only course runs starting strictly before this instant."), + ] = None, ) -> LearnerRecordsResponse[BaseModel]: - filters = queries.RecordFilters(organization_id=organization_id, contract_id=contract_id) + filters = queries.RecordFilters( + organization_id=organization_id, + contract_id=contract_id, + contract_is_active=contract_is_active, + courserun_id=courserun_id, + courserun_starts_after=courserun_starts_after, + courserun_starts_before=courserun_starts_before, + ) return await _respond( queries.courses(settings.starrocks_schema, filters), organization_id, page, CourseRun ) diff --git a/tests/test_learner_records.py b/tests/test_learner_records.py index 4cd09be..16c2c06 100644 --- a/tests/test_learner_records.py +++ b/tests/test_learner_records.py @@ -480,6 +480,22 @@ async def test_enrollment_filters_are_bound_in_order(app, monkeypatch): assert pool.count_call()[1] == params[:-2] +async def test_updated_before_bounds_the_window_exclusively(app, monkeypatch): + pool = _FakePool() + await _get( + app, + f"/organizations/{ORG_ID}/enrollments" + "?updated_since=2026-08-12T00:00:00Z&updated_before=2026-08-13T00:00:00Z", + _partner_header(ORG_ID), + pool, + monkeypatch, + ) + query, params = pool.page_call() + assert "record_updated_on >= %s" in query + assert "record_updated_on < %s" in query + assert params == (ORG_ID, "2026-08-12T00:00:00", "2026-08-13T00:00:00", 100, 0) + + async def test_omitted_list_filters_add_no_predicate(app, monkeypatch): """`learner_id` and `completion_status` are plain arrays defaulting to [], not `list | None`: the optional form publishes `anyOf: [array, null]`, which @@ -532,6 +548,25 @@ def test_cursor_value_is_a_prefix_of_stored_values_in_the_same_second(): assert stored >= bound +def test_round_up_cursor_value_never_excludes_a_row_before_the_requested_instant(): + # A row stored earlier in the same second as a fractional exclusive bound + # must still satisfy `< bound`. Flooring 06:15:00.500 to "...:00" would put + # a row stored at "...:00.250" (which precedes .500) on the wrong side of + # `<`, silently dropping it rather than re-sending it in the next window. + bound = queries._cursor_value( # noqa: SLF001 + datetime.datetime(2026, 8, 12, 6, 15, 0, 500000, tzinfo=datetime.UTC), round_up=True + ) + assert bound == "2026-08-12T06:15:01" + assert bound > "2026-08-12T06:15:00.250000" + + +def test_round_up_cursor_value_is_unchanged_on_a_whole_second(): + bound = queries._cursor_value( # noqa: SLF001 + datetime.datetime(2026, 8, 12, 6, 15, 0, tzinfo=datetime.UTC), round_up=True + ) + assert bound == "2026-08-12T06:15:00" + + def test_tenant_never_imports_the_anonymization_module(): """The two tenants' privacy postures should be legible from their imports.""" package = pathlib.Path(b2b_learner_records.__file__).parent @@ -623,3 +658,31 @@ async def test_courses_contract_filter_is_bound(app, monkeypatch): assert "sso_organization_id = %s AND contract_id = %s" in query assert params == (ORG_ID, 42, 100, 0) assert pool.count_call()[1] == (ORG_ID, 42) + + +async def test_courses_filters_are_bound_in_order(app, monkeypatch): + pool = _FakePool() + await _get( + app, + f"/organizations/{ORG_ID}/courses" + "?contract_is_active=true&courserun_id=course-v1:MITxT%2B14.310x%2B2T2026" + "&courserun_starts_after=2026-01-01T00:00:00Z&courserun_starts_before=2026-04-01T00:00:00Z", + _partner_header(ORG_ID), + pool, + monkeypatch, + ) + query, params = pool.page_call() + assert "b2b_contract_is_active = %s" in query + assert "courserun_readable_id = %s" in query + assert "courserun_start_on >= %s" in query + assert "courserun_start_on < %s" in query + assert params == ( + ORG_ID, + True, + "course-v1:MITxT+14.310x+2T2026", + "2026-01-01T00:00:00", + "2026-04-01T00:00:00", + 100, + 0, + ) + assert pool.count_call()[1] == params[:-2]