Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
70 changes: 70 additions & 0 deletions openapi/specs/b2b_learner_records.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
83 changes: 72 additions & 11 deletions src/ol_analytics_api/tenants/b2b_learner_records/queries.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand All @@ -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
Expand Down Expand Up @@ -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


Expand Down Expand Up @@ -384,13 +426,32 @@ 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"]
params: list[Any] = [str(filters.organization_id)]
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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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).")
]
Expand Down Expand Up @@ -109,13 +118,15 @@ 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(
organization_id=organization_id,
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(
Expand All @@ -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(
Expand All @@ -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(
Expand All @@ -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
)
63 changes: 63 additions & 0 deletions tests/test_learner_records.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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]
Loading