Skip to content
Draft
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
120 changes: 118 additions & 2 deletions openedx/core/djangoapps/content/search/api.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

from __future__ import annotations

import json
import logging
import time
from collections.abc import Iterator
Expand All @@ -16,11 +17,12 @@
from django.conf import settings
from django.contrib.auth import get_user_model
from django.core.cache import cache
from django.core.exceptions import ObjectDoesNotExist
from django.core.paginator import Paginator
from meilisearch import Client as MeilisearchClient
from meilisearch.errors import MeilisearchApiError, MeilisearchError
from meilisearch.models.task import TaskInfo
from opaque_keys import OpaqueKey
from opaque_keys import InvalidKeyError, OpaqueKey
from opaque_keys.edx.keys import CourseKey, LearningContextKey, UsageKey
from opaque_keys.edx.locator import LibraryCollectionLocator, LibraryContainerLocator, LibraryLocatorV2
from openedx_content import api as content_api
Expand All @@ -38,11 +40,16 @@
INDEX_SEARCHABLE_ATTRIBUTES,
INDEX_SORTABLE_ATTRIBUTES,
)
from openedx.core.djangoapps.content.search.models import IncrementalIndexCompleted, get_access_ids_for_request
from openedx.core.djangoapps.content.search.models import (
IncrementalIndexCompleted,
SearchAccess,
get_access_ids_for_request,
)
from openedx.core.djangoapps.content_libraries import api as lib_api
from xmodule.modulestore.django import modulestore
from xmodule.modulestore.exceptions import ItemNotFoundError

from .content_reconciliation import reconcile_components, validate_reconciliation_limits
from .documents import (
DocType,
Fields,
Expand Down Expand Up @@ -1306,3 +1313,112 @@ def get_all_blocks_from_context(
break

offset += limit


def reconcile_library_components(
library_key: LibraryLocatorV2, *, repair: bool = False,
batch_size: int = 100, max_documents: int = 10000,
):
"""Inspect/repair one library's component documents without resetting an index.

Containers/collections and index-only documents are not mutated. Call from a
management command, never a request handler (Meilisearch tasks are awaited).
Content publication must be quiesced for strict repair correctness; source
and index rereads reduce races but cannot provide cross-system atomicity.
"""
if not isinstance(library_key, LibraryLocatorV2):
raise ValueError("A single Libraries V2 key is required")
validate_reconciliation_limits(batch_size, max_documents)
# Validate the source library before any engine requests, and require the
# pre-existing access row so dry-run serialization doesn't create one.
lib_api.get_library(library_key)
if not SearchAccess.objects.filter(context_key=str(library_key)).exists():
raise ValueError("Library search access metadata is missing; populate the index first")
client = _get_meilisearch_client()
index = client.get_index(STUDIO_LIBRARY_INDEX_NAME)
if index.primary_key != INDEX_PRIMARY_KEY:
raise ValueError("Reconcile index settings before reconciling content")

def check_rebuild():
if _get_running_rebuild_index_name(STUDIO_LIBRARY_INDEX_NAME):
raise RuntimeError("Library index rebuild in progress; retry reconciliation afterwards")

check_rebuild()

def source_keys():
components = lib_api.get_library_components(library_key).order_by("pk")
for component in components.iterator(chunk_size=batch_size):
yield lib_api.LibraryXBlockMetadata.from_component(library_key, component).usage_key

def read_source(key):
key = UsageKey.from_string(str(key))
if key.context_key != library_key:
raise ValueError("Component does not belong to selected library")
try:
component = lib_api.get_component_from_usage_key(key)
# Published-only entities are deliberately outside this draft index.
if not lib_api.get_library_components(library_key).filter(pk=component.pk).exists():
return None
except (ObjectDoesNotExist, lib_api.ContentLibraryBlockNotFound):
return None
metadata = lib_api.LibraryXBlockMetadata.from_component(library_key, component)
doc = searchable_doc_for_library_block(metadata)
doc.update(searchable_doc_tags(key))
doc.update(searchable_doc_collections(key))
doc.update(searchable_doc_containers(key, "units"))
return doc

def as_dict(document):
# The Meilisearch SDK returns Document instances with dynamic fields.
return dict(document) if isinstance(document, dict) else vars(document)

def read_index(document_id):
try:
return as_dict(index.get_document(document_id))
except MeilisearchApiError as err:
if err.code == "document_not_found":
return None
raise

def indexed_documents():
offset = 0
# Documents endpoint avoids search maxTotalHits and distinct-attribute
# truncation. JSON quoting prevents filter interpolation.
while offset <= max_documents:
response = index.get_documents({
"filter": (
f"context_key = {json.dumps(str(library_key))} "
f"AND type = {json.dumps(DocType.library_block)}"
),
"offset": offset,
"limit": min(batch_size, max_documents + 1 - offset),
})
if not response.results:
return
for document in response.results:
doc = as_dict(document)
# Corrupt indexed keys are unknown, not evidence of deletion.
# Only key parsing errors are swallowed; source/engine failures
# still abort the audit or repair.
try:
value = doc.get(Fields.usage_key)
key = UsageKey.from_string(value) if isinstance(value, str) else None
except (InvalidKeyError, ValueError):
key = None
if key is None or key.context_key != library_key:
doc = {**doc, Fields.usage_key: None}
yield doc
offset += len(response.results)

def write_batch(docs):
check_rebuild()
# add_documents replaces entire records; update_documents would retain
# obsolete fields and leave drift behind. No deletion tasks are issued.
_wait_for_meili_task(index.add_documents(docs))
check_rebuild()

return reconcile_components(
context_key=str(library_key), source_keys=source_keys(), read_source=read_source,
read_index=read_index, indexed_documents=indexed_documents(), write_batch=write_batch,
repair=repair, batch_size=batch_size, max_documents=max_documents,
)
129 changes: 129 additions & 0 deletions openedx/core/djangoapps/content/search/content_reconciliation.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
"""Bounded, source-authoritative library component reconciliation.

This module deliberately has no Django imports so its repair policy can be tested
without a platform runtime. The platform adapters live in ``api.py``.
"""

from dataclasses import asdict, dataclass
from itertools import islice


@dataclass
class ContentDriftReport:
"""Counters only: reports never expose component bodies or access identifiers."""

scanned: int = 0
missing: int = 0
stale: int = 0
unchanged: int = 0
repaired: int = 0
skipped_concurrent: int = 0
index_only: int = 0
unknown: int = 0
source_truncated: bool = False
index_truncated: bool = False

def as_dict(self):
"""Return a JSON-serializable summary."""
return asdict(self)


def validate_reconciliation_limits(batch_size, max_documents):
"""Require integer scan/batch limits and reject booleans explicitly."""
if not isinstance(batch_size, int) or isinstance(batch_size, bool) or not 1 <= batch_size <= 1000:
raise ValueError("batch_size must be between 1 and 1000")
if not isinstance(max_documents, int) or isinstance(max_documents, bool) or max_documents < 1:
raise ValueError("max_documents must be positive")


def _valid_component_identity(doc, context_key):
"""Require a component identity scoped to the selected library."""
return (
doc is not None and doc.get("context_key") == context_key
and doc.get("type") == "library_block" and bool(doc.get("id"))
and bool(doc.get("usage_key"))
)


def _inspect_indexed_components(report, *, context_key, indexed_documents, read_source, max_documents):
"""Inspect bounded indexed identities without authorizing any deletion."""
for position, doc in enumerate(islice(indexed_documents, max_documents + 1)):
if position == max_documents:
report.index_truncated = True
break
if not _valid_component_identity(doc, context_key):
report.unknown += 1
else:
canonical = read_source(doc["usage_key"])
if canonical is None:
report.index_only += 1
elif canonical["id"] != doc["id"]:
report.unknown += 1


def reconcile_components(
*, context_key, source_keys, read_source, read_index, indexed_documents,
write_batch, repair=False, batch_size=100, max_documents=10000,
):
"""Compare canonical source documents and cautiously replace drifted entries.

``read_source`` must return None for missing/deleted draft components and
propagate all other failures. ``read_index`` likewise treats only an explicit
document-not-found result as absence. Source keys and index documents must be
lazy bounded streams. Full-document equality detects removed fields as drift.
Index-only documents are reported, never deleted. Each scan has its own cap.
"""
validate_reconciliation_limits(batch_size, max_documents)
report = ContentDriftReport()
pending = []

def flush():
docs = []
for key, expected, observed in pending:
# Revalidate immediately before the bounded write, not at scan time.
latest = read_source(key)
current = read_index(expected["id"])
if latest != expected or current != observed:
report.skipped_concurrent += 1
else:
docs.append(latest)
if docs:
write_batch(docs) # Must wait for task success; exceptions abort the run.
report.repaired += len(docs)
pending.clear()

for position, key in enumerate(islice(source_keys, max_documents + 1)):
if position == max_documents:
report.source_truncated = True
break
report.scanned += 1
expected = read_source(key)
if not _valid_component_identity(expected, context_key) or expected["usage_key"] != str(key):
report.unknown += 1
continue
observed = read_index(expected["id"])
# A colliding/corrupt primary key must not replace another context/type.
if observed is not None and (
not _valid_component_identity(observed, context_key) or observed["usage_key"] != expected["usage_key"]
):
report.unknown += 1
continue
if expected == observed:
report.unchanged += 1
continue
if observed is None:
report.missing += 1
else:
report.stale += 1
if repair:
pending.append((key, expected, observed))
if len(pending) == batch_size:
flush()
if pending:
flush()

_inspect_indexed_components(
report, context_key=context_key, indexed_documents=indexed_documents,
read_source=read_source, max_documents=max_documents,
)
return report
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
Library component drift audit and targeted repair
=================================================

Problem and existing behavior
-----------------------------

``reconcile_indexes`` handles index schema/settings; ``reindex_studio`` queues
population or an index rebuild. Neither compares individual live component
records with canonical source serialization. This contribution adds an explicit
operator command and API for one Libraries V2 context. It does not replace the
existing reconciliation or reindex pathways.

Implementation
--------------

``reconcile_library_components`` validates the library and existing search-access
metadata, requires the existing library index/primary key, then calls the
stdlib-only ``reconcile_components`` policy. Source enumeration uses the existing
draft-components queryset in primary-key order and ``iterator(chunk_size=...)``.
The document is assembled with the existing library-block serializer, tags,
collections and units, matching the full rebuild fields including permissions.
Each source document is compared against ``get_document``. Only the SDK's
``document_not_found`` error is absence; missing indexes, invalid credentials,
timeouts and all other service failures abort.

Source and index scans each have an explicit cap (default 10,000) with separate
truncation indicators. One lookahead item detects truncation. Index inspection
uses the SDK documents endpoint, filtered to exactly the selected context and
``library_block`` type. Numeric offsets advance by returned page lengths, with
requests bounded by batch size and remaining lookahead budget. This avoids
search-result limits and distinct grouping, but offset pagination is not a
snapshot. No all-library list or whole-library document buffer is built.

Dry-run is the default. The report contains counts, never content bodies or
permission identifiers. A library must already have SearchAccess metadata so
normal serialization's get-or-create path finds an existing row. This is not a
database read-only transaction; concurrently deleting that access row may still
cause the serializer to recreate it. Operators should keep metadata stable.

Repair policy
-------------

* An explicit ``--repair`` is required for search writes.
* Missing and stale canonical component records are eligible. Unknown/corrupt
identities and primary-key collisions with another context/type/key are skipped.
* Source and index records are both re-read immediately before each bounded batch.
If either differs from the observed candidate, repair skips it as concurrent.
* ``add_documents`` replaces the complete record, including removal of obsolete
fields; partial ``update_documents`` would leave such drift behind.
* Each submitted task must succeed before its batch is counted as repaired.
Failure aborts; earlier successful batches remain applied and a later run can
inspect them again.
* An active library-index rebuild is refused before scanning and before/after
each write. The course index and temporary indexes are never written.
* Index-only records are independently checked against source, reported, and
never deleted. Source-enumeration truncation cannot create false orphans because
indexed keys are checked directly against source existence.
* Corrupt or wrong-context indexed usage keys count as unknown. Only key parsing
failures are absorbed; related content/engine failures propagate.

Concurrency and completeness limits
-----------------------------------

This is a cautious repair tool, not an atomic cross-system reconciler. Component,
association, tags, permissions and index records do not share a transaction or
compare-and-swap primitive. A publication after the final reread, a previously
queued newer indexing task, or an index swap between the rebuild check and write
can still race. Strict repair correctness requires pausing content mutations,
draining indexing tasks and avoiding rebuilds during the operation. A repair
run should be followed by a dry-run after writes have settled. Offset pagination
may miss/revisit documents while additions occur; it never authorizes deletion.

A truncated report is not a complete-library health assertion. Increase
``--max-documents`` above the selected library's size and rerun under quiescence.
Persistent cursor/keyset support is future work. Collection and container repair
are deferred, because their deletion/draft semantics differ from components.
Automatic orphan deletion, scheduler/Celery integration, distributed fencing,
revision-aware conditional writes and repair audit persistence are also deferred.

Validation and contribution sequencing
--------------------------------------

The dependency-free policy suite tests batching, truncation, dry-run, no-delete,
identity collisions, removed fields, failures and source/index mutation races.
CMS integration tests create real libraries, components, a collection and tags,
and exercise Django command discovery and canonical serializers. Only the
external engine boundary is replaced in those cases. They cover missing access
metadata, invalid/missing library keys, rebuild refusal, scan caps, wrong primary
keys, corrupt identities, collisions and engine failures.

An opt-in live-engine test uses the same real CMS fixtures and a unique disposable
index. Fixture creation disables indexing events; content stays unchanged after
initialization and all engine tasks are awaited. It checks a missing document,
a stale payload with an obsolete field, an index-only record and another
library. Dry-run performs no writes, repair restores both canonical records,
a second audit is clean, and the unrelated records are preserved. The temporary
index is removed in a finally block.

Focused validation passed against the platform master base with Meilisearch
1.36.0: 18 policy cases, 14 CMS integration cases, one live-engine case,
34 existing index-settings reconciliation cases, and the library creation and
component creation/deletion handler cases (69 distinct tests). After bounds and indexed-identity fixes, only the
33 affected policy/CMS/live cases were rerun.
Repository-configured Ruff and Pylint, and git diff --check, passed.

Run the targeted tests using the CMS test settings. The live test is skipped
unless OPENEDX_SEARCH_LIVE_URL is supplied. Set OPENEDX_SEARCH_LIVE_KEY privately
to an engine key that can create, update and delete disposable indexes; do not
use a production service. No frontend, LMS or load suite is needed for this
operator-only contribution.

Maintainer review should confirm publication quiescence procedures. Collection
and container repair and revision-fenced background operation should be proposed
separately. These tests establish controlled repair behavior, not production
concurrency guarantees or performance claims.
Loading