From c05137420bdf3651c5b48c5ec2e9786528eed18f Mon Sep 17 00:00:00 2001 From: Shatrughna Upadhyay Date: Sat, 3 Oct 2026 13:37:57 -0700 Subject: [PATCH 1/2] feat: add bounded library component search reconciliation --- openedx/core/djangoapps/content/search/api.py | 119 ++++++++++- .../content/search/content_reconciliation.py | 110 ++++++++++ .../docs/library-content-reconciliation.rst | 95 +++++++++ .../commands/reconcile_library_search.py | 34 +++ .../test_content_reconciliation_adapter.py | 201 ++++++++++++++++++ .../test_content_reconciliation_policy.py | 142 +++++++++++++ 6 files changed, 700 insertions(+), 1 deletion(-) create mode 100644 openedx/core/djangoapps/content/search/content_reconciliation.py create mode 100644 openedx/core/djangoapps/content/search/docs/library-content-reconciliation.rst create mode 100644 openedx/core/djangoapps/content/search/management/commands/reconcile_library_search.py create mode 100644 openedx/core/djangoapps/content/search/tests/test_content_reconciliation_adapter.py create mode 100644 openedx/core/djangoapps/content/search/tests/test_content_reconciliation_policy.py diff --git a/openedx/core/djangoapps/content/search/api.py b/openedx/core/djangoapps/content/search/api.py index da68dcacc297..e4ef89c68564 100644 --- a/openedx/core/djangoapps/content/search/api.py +++ b/openedx/core/djangoapps/content/search/api.py @@ -20,7 +20,7 @@ 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 @@ -1306,3 +1306,120 @@ 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. + """ + import json # pylint: disable=import-outside-toplevel + + from django.core.exceptions import ObjectDoesNotExist # pylint: disable=import-outside-toplevel + + from .content_reconciliation import reconcile_components # pylint: disable=import-outside-toplevel + from .models import SearchAccess # pylint: disable=import-outside-toplevel + + if not isinstance(library_key, LibraryLocatorV2): + raise ValueError("A single Libraries V2 key is required") + if ( + type(batch_size) is not int or not 1 <= batch_size <= 1000 + or type(max_documents) is not int or max_documents < 1 + ): + raise ValueError("Invalid batch_size or 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))} 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, + ) diff --git a/openedx/core/djangoapps/content/search/content_reconciliation.py b/openedx/core/djangoapps/content/search/content_reconciliation.py new file mode 100644 index 000000000000..069a94a86378 --- /dev/null +++ b/openedx/core/djangoapps/content/search/content_reconciliation.py @@ -0,0 +1,110 @@ +"""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 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. + """ + if type(batch_size) is not int or not 1 <= batch_size <= 1000: + raise ValueError("batch_size must be between 1 and 1000") + if type(max_documents) is not int or max_documents < 1: + raise ValueError("max_documents must be positive") + report = ContentDriftReport() + pending = [] + + def valid(doc): + 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 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(expected) 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(observed) 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() + + for position, doc in enumerate(islice(indexed_documents, max_documents + 1)): + if position == max_documents: + report.index_truncated = True + break + if not valid(doc): + report.unknown += 1 + elif read_source(doc["usage_key"]) is None: + report.index_only += 1 + return report diff --git a/openedx/core/djangoapps/content/search/docs/library-content-reconciliation.rst b/openedx/core/djangoapps/content/search/docs/library-content-reconciliation.rst new file mode 100644 index 000000000000..db343062925c --- /dev/null +++ b/openedx/core/djangoapps/content/search/docs/library-content-reconciliation.rst @@ -0,0 +1,95 @@ +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. +The isolated API adapter suite executes the actual function body with mocked +platform services and the pinned Meilisearch 0.43.0 Document/exception classes. +It checks filter scoping, page offsets/limits, full replacement, permission fields, +missing-document versus missing-index errors, rebuild checks and access metadata. +These are not full Django/database or live-engine integration tests. + +Before proposing an upstream PR, add Django library fixture tests and run the +existing search suite with a real Meilisearch. Reproduce a missed indexing event, +verify repaired content through Studio, and confirm publication quiescence +procedures with maintainers. Start with this bounded operator-only scope; propose +container/collection repair and revision-fenced background operation separately. diff --git a/openedx/core/djangoapps/content/search/management/commands/reconcile_library_search.py b/openedx/core/djangoapps/content/search/management/commands/reconcile_library_search.py new file mode 100644 index 000000000000..2802a3c3a42e --- /dev/null +++ b/openedx/core/djangoapps/content/search/management/commands/reconcile_library_search.py @@ -0,0 +1,34 @@ +"""Inspect and optionally repair component drift in one library.""" + +import json + +from django.core.management import BaseCommand, CommandError +from opaque_keys import InvalidKeyError +from opaque_keys.edx.locator import LibraryLocatorV2 + +from ... import api + + +class Command(BaseCommand): + """Default to inspection, requiring explicit --repair for document writes.""" + + help = "Inspect component search drift in one library; --repair replaces missing/stale documents." + + def add_arguments(self, parser): + parser.add_argument("library_key", help="Libraries V2 key, e.g. lib:org:slug") + parser.add_argument("--repair", action="store_true", help="Write repairs (default: dry run)") + parser.add_argument("--batch-size", type=int, default=100) + parser.add_argument("--max-documents", type=int, default=10000, help="Cap each source/index scan") + + def handle(self, *args, **options): + if not api.is_meilisearch_enabled(): + raise CommandError("Meilisearch is disabled") + try: + key = LibraryLocatorV2.from_string(options["library_key"]) + report = api.reconcile_library_components( + key, repair=options["repair"], batch_size=options["batch_size"], + max_documents=options["max_documents"], + ) + except (InvalidKeyError, ValueError, RuntimeError) as err: + raise CommandError(str(err)) from err + self.stdout.write(json.dumps({"repair": options["repair"], **report.as_dict()}, sort_keys=True)) diff --git a/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_adapter.py b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_adapter.py new file mode 100644 index 000000000000..d0c6f28fba69 --- /dev/null +++ b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_adapter.py @@ -0,0 +1,201 @@ +"""Isolated adapter contract tests with the pinned Meilisearch SDK. + +Execute the real API function body with mocked platform services. This validates +adapter behavior, not Django ORM/database integration or a live Meilisearch. +""" +# ruff: noqa: PT009, PT027 +import ast +import importlib.util +import sys +import types +import unittest +from pathlib import Path +from unittest.mock import MagicMock, patch + +from meilisearch.errors import MeilisearchApiError +from meilisearch.models.document import Document, DocumentsResults + +ROOT = Path(__file__).parents[1] +SPEC = importlib.util.spec_from_file_location("isolated_search.content_reconciliation", ROOT / "content_reconciliation.py") +POLICY = importlib.util.module_from_spec(SPEC) +sys.modules[SPEC.name] = POLICY +SPEC.loader.exec_module(POLICY) + + +class LibraryKey: + """Minimal stand-in for an already-parsed library key.""" + + def __str__(self): + return "lib:org:one" + + +class ComponentKey: + """Minimal usage key satisfying the API's parsing/context contract.""" + + def __init__(self, name): + self.name = name + self.context_key = LIBRARY + + def __str__(self): + return self.name + + @classmethod + def from_string(cls, key): + return cls(key) + + +LIBRARY = LibraryKey() + + +def doc(key="a", **fields): + return {"id": key, "usage_key": key, "context_key": str(LIBRARY), "type": "library_block", **fields} + + +class NotFound(Exception): + """Stand-in for Django ObjectDoesNotExist.""" + + +def engine_error(code): + # Use the real SDK exception class, without constructing an HTTP response. + err = MeilisearchApiError.__new__(MeilisearchApiError) + Exception.__init__(err, code) + err.code = code + return err + + +class AdapterTest(unittest.TestCase): + """Verify the actual API adapter, with only external dependencies replaced.""" + + def setUp(self): + self.lib = MagicMock() + self.index = MagicMock(primary_key="id") + self.client = MagicMock() + self.client.get_index.return_value = self.index + self.access = MagicMock() + self.access.objects.filter.return_value.exists.return_value = True + self.lib.get_library_components.return_value.order_by.return_value.iterator.return_value = ["a"] + self.lib.get_library_components.return_value.filter.return_value.exists.return_value = True + self.lib.get_component_from_usage_key.side_effect = lambda key: types.SimpleNamespace(pk=str(key)) + self.lib.LibraryXBlockMetadata.from_component.side_effect = lambda key, comp: types.SimpleNamespace( + usage_key=ComponentKey(comp if isinstance(comp, str) else comp.pk), + ) + self.lib.ContentLibraryBlockNotFound = NotFound + self.index.get_document.return_value = Document(doc()) + self.index.get_documents.return_value = DocumentsResults({"results": [], "offset": 0, "limit": 100, "total": 0}) + self.rebuild = MagicMock(return_value=None) + self.wait = MagicMock() + self.namespace = { + "__package__": "isolated_search", "LibraryLocatorV2": LibraryKey, "UsageKey": ComponentKey, + "lib_api": self.lib, "STUDIO_LIBRARY_INDEX_NAME": "library", "INDEX_PRIMARY_KEY": "id", + "DocType": types.SimpleNamespace(library_block="library_block"), + "Fields": types.SimpleNamespace(usage_key="usage_key"), "InvalidKeyError": ValueError, + "MeilisearchApiError": MeilisearchApiError, "_get_meilisearch_client": lambda: self.client, + "_get_running_rebuild_index_name": self.rebuild, "_wait_for_meili_task": self.wait, + "searchable_doc_for_library_block": lambda metadata: doc(str(metadata.usage_key), access_id=7), + "searchable_doc_tags": lambda key: {"tags": {}}, + "searchable_doc_collections": lambda key: {"collections": {}}, + "searchable_doc_containers": lambda key, group: {group: {}}, + } + tree = ast.parse((ROOT / "api.py").read_text()) + target = next(node for node in tree.body if isinstance(node, ast.FunctionDef) + and node.name == "reconcile_library_components") + exec(compile(ast.Module(body=[target], type_ignores=[]), "api.py", "exec"), self.namespace) + self.function = self.namespace["reconcile_library_components"] + exceptions = types.ModuleType("django.core.exceptions") + exceptions.ObjectDoesNotExist = NotFound + models = types.ModuleType("isolated_search.models") + models.SearchAccess = self.access + self.modules = patch.dict(sys.modules, {"django.core.exceptions": exceptions, "isolated_search.models": models}) + self.modules.start() + self.addCleanup(self.modules.stop) + + def test_dry_run_document_sdk_and_no_writes(self): + report = self.function(LIBRARY) + self.assertEqual(report.stale, 1) + self.index.add_documents.assert_not_called() + self.access.objects.create.assert_not_called() + self.assertEqual(self.index.get_documents.call_args.args[0]["filter"], + 'context_key = "lib:org:one" AND type = "library_block"') + + def test_repair_full_document_and_wait(self): + report = self.function(LIBRARY, repair=True) + self.assertEqual(report.repaired, 1) + written = self.index.add_documents.call_args.args[0][0] + self.assertEqual(written["access_id"], 7) + self.assertEqual(written["tags"], {}) + self.assertEqual(written["units"], {}) + self.index.update_documents.assert_not_called() + self.wait.assert_called_once_with(self.index.add_documents.return_value) + + def test_only_document_not_found_is_absence(self): + self.index.get_document.side_effect = engine_error("document_not_found") + self.assertEqual(self.function(LIBRARY).missing, 1) + + def test_missing_index_not_treated_as_missing_document(self): + self.index.get_document.side_effect = engine_error("index_not_found") + with self.assertRaises(MeilisearchApiError): + self.function(LIBRARY, repair=True) + self.index.add_documents.assert_not_called() + + def test_engine_authorization_failure_aborts(self): + self.index.get_document.side_effect = engine_error("invalid_api_key") + with self.assertRaises(MeilisearchApiError): + self.function(LIBRARY, repair=True) + self.index.add_documents.assert_not_called() + + def test_active_rebuild_refused(self): + self.rebuild.return_value = "library_new" + with self.assertRaises(RuntimeError): + self.function(LIBRARY, repair=True) + self.index.add_documents.assert_not_called() + + def test_rebuild_started_before_write_refused(self): + self.rebuild.side_effect = [None, "library_new"] + with self.assertRaises(RuntimeError): + self.function(LIBRARY, repair=True) + self.index.add_documents.assert_not_called() + + def test_missing_access_dry_run_refused_without_creation(self): + self.access.objects.filter.return_value.exists.return_value = False + with self.assertRaises(ValueError): + self.function(LIBRARY) + self.client.get_index.assert_not_called() + + def test_paginated_documents_bounded_offsets(self): + self.lib.get_library_components.return_value.order_by.return_value.iterator.return_value = [] + self.index.get_documents.side_effect = [ + DocumentsResults({"results": [doc("a"), doc("b")], "offset": 0, "limit": 2, "total": 4}), + DocumentsResults({"results": [doc("c")], "offset": 2, "limit": 1, "total": 4}), + ] + report = self.function(LIBRARY, batch_size=2, max_documents=2) + self.assertTrue(report.index_truncated) + self.assertEqual([call.args[0]["offset"] for call in self.index.get_documents.call_args_list], [0, 2]) + self.assertEqual([call.args[0]["limit"] for call in self.index.get_documents.call_args_list], [2, 1]) + + def test_malformed_index_key_reported_unknown(self): + self.lib.get_library_components.return_value.order_by.return_value.iterator.return_value = [] + self.index.get_documents.side_effect = [ + DocumentsResults({"results": [doc(usage_key=42)], "offset": 0, "limit": 100, "total": 1}), + DocumentsResults({"results": [], "offset": 1, "limit": 100, "total": 1}), + ] + report = self.function(LIBRARY) + self.assertEqual(report.unknown, 1) + self.assertEqual(report.index_only, 0) + + def test_wrong_primary_key_refused(self): + self.index.primary_key = "other" + with self.assertRaises(ValueError): + self.function(LIBRARY, repair=True) + self.index.add_documents.assert_not_called() + + def test_builder_failure_propagates_not_orphan(self): + def fail_builder(_): + raise NotFound("related source read failed") + self.namespace["searchable_doc_for_library_block"] = fail_builder + with self.assertRaises(NotFound): + self.function(LIBRARY, repair=True) + self.index.add_documents.assert_not_called() + + +if __name__ == "__main__": + unittest.main() diff --git a/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_policy.py b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_policy.py new file mode 100644 index 000000000000..463e4db39004 --- /dev/null +++ b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_policy.py @@ -0,0 +1,142 @@ +"""Dependency-free policy tests; runnable with unittest outside Django settings.""" + +# unittest is intentional: this policy suite runs without pytest/Django installed. +# ruff: noqa: PT009, PT027 +import importlib.util +import sys +import unittest +from pathlib import Path + +# Load just the policy file, not the platform's Django application graph. +_SPEC = importlib.util.spec_from_file_location( + "content_reconciliation_policy", Path(__file__).parents[1] / "content_reconciliation.py", +) +_POLICY = importlib.util.module_from_spec(_SPEC) +sys.modules[_SPEC.name] = _POLICY +_SPEC.loader.exec_module(_POLICY) +reconcile_components = _POLICY.reconcile_components + + +def document(key="a", **fields): + return {"id": key, "usage_key": key, "context_key": "lib:org:one", "type": "library_block", **fields} + + +class ReconcilePolicyTest(unittest.TestCase): + """Exercise bounded writes, ambiguous identity and observed mutation races.""" + + def run_reconcile(self, sources=None, indexed=None, **kwargs): + sources = {"a": document()} if sources is None else sources + indexed = {} if indexed is None else indexed + self.writes = [] + + def write_batch(docs): + self.writes.append(docs) + indexed.update({doc["id"]: doc for doc in docs}) + + options = { + "context_key": "lib:org:one", "source_keys": iter(sources), + "read_source": sources.get, "read_index": indexed.get, + "indexed_documents": iter(list(indexed.values())), "write_batch": write_batch, + } + options.update(kwargs) + return reconcile_components(**options) + + def test_dry_run_default(self): + report = self.run_reconcile() + self.assertEqual(report.missing, 1) + self.assertEqual(report.repaired, 0) + self.assertEqual(self.writes, []) + + def test_missing_repair(self): + report = self.run_reconcile(repair=True) + self.assertEqual(report.repaired, 1) + self.assertEqual(self.writes, [[document()]]) + + def test_stale_removed_field_replaced(self): + report = self.run_reconcile(indexed={"a": document(obsolete="remove me")}, repair=True) + self.assertEqual((report.stale, report.repaired), (1, 1)) + self.assertNotIn("obsolete", self.writes[0][0]) + + def test_unchanged_no_write(self): + report = self.run_reconcile(indexed={"a": document()}, repair=True) + self.assertEqual(report.unchanged, 1) + self.assertEqual(self.writes, []) + + def test_batch_bound(self): + report = self.run_reconcile(sources={str(n): document(str(n)) for n in range(7)}, repair=True, batch_size=3) + self.assertEqual(report.repaired, 7) + self.assertEqual([len(batch) for batch in self.writes], [3, 3, 1]) + + def test_source_bound_explicit(self): + report = self.run_reconcile(sources={str(n): document(str(n)) for n in range(7)}, max_documents=2) + self.assertEqual(report.scanned, 2) + self.assertTrue(report.source_truncated) + + def test_index_bound_explicit(self): + report = self.run_reconcile(sources={}, indexed={str(n): document(str(n)) for n in range(7)}, max_documents=2) + self.assertEqual(report.index_only, 2) + self.assertTrue(report.index_truncated) + + def test_index_only_never_deleted(self): + report = self.run_reconcile(sources={}, indexed={"a": document()}, repair=True) + self.assertEqual(report.index_only, 1) + self.assertEqual(self.writes, []) + + def test_unenumerated_source_is_not_orphan(self): + report = self.run_reconcile( + sources={"a": document(), "b": document("b")}, + indexed={"b": document("b")}, max_documents=1, + ) + self.assertTrue(report.source_truncated) + self.assertEqual(report.index_only, 0) + + def test_foreign_context_collision_skipped(self): + report = self.run_reconcile(indexed={"a": document(context_key="lib:org:other")}, repair=True) + self.assertEqual(report.unknown, 2) + self.assertEqual(self.writes, []) + + def test_wrong_key_collision_skipped(self): + report = self.run_reconcile(indexed={"a": document(usage_key="other")}, repair=True) + self.assertGreaterEqual(report.unknown, 1) + self.assertEqual(self.writes, []) + + def test_source_change_before_write_skipped(self): + reads = iter([document(), document(display_name="new")]) + report = self.run_reconcile(repair=True, read_source=lambda _: next(reads)) + self.assertEqual(report.skipped_concurrent, 1) + self.assertEqual(self.writes, []) + + def test_source_deleted_before_write_skipped(self): + reads = iter([document(), None]) + report = self.run_reconcile(repair=True, read_source=lambda _: next(reads)) + self.assertEqual(report.skipped_concurrent, 1) + self.assertEqual(self.writes, []) + + def test_index_change_before_write_skipped(self): + reads = iter([None, document(display_name="new")]) + report = self.run_reconcile(repair=True, read_index=lambda _: next(reads)) + self.assertEqual(report.skipped_concurrent, 1) + self.assertEqual(self.writes, []) + + def test_transport_error_not_missing(self): + def failed_read(_): + raise ConnectionError("engine unavailable") + with self.assertRaises(ConnectionError): + self.run_reconcile(read_index=failed_read, repair=True) + self.assertEqual(self.writes, []) + + def test_failed_write_propagates(self): + def failed_write(_): + raise RuntimeError("task failed") + with self.assertRaises(RuntimeError): + self.run_reconcile(repair=True, write_batch=failed_write) + + def test_invalid_limits(self): + for kwargs in ({"batch_size": 0}, {"batch_size": 1001}, {"max_documents": 0}, + {"batch_size": True}, {"batch_size": 1.5}, {"max_documents": True}): + with self.subTest(kwargs=kwargs), self.assertRaises(ValueError): + self.run_reconcile(**kwargs) + + +if __name__ == "__main__": + unittest.main() From f8275c49cc2ed582975c3b3e233848b507ff868f Mon Sep 17 00:00:00 2001 From: Shatrughna Upadhyay Date: Sat, 3 Oct 2026 14:44:55 -0700 Subject: [PATCH 2/2] test: validate library reconciliation in CMS and Meilisearch --- openedx/core/djangoapps/content/search/api.py | 27 +- .../content/search/content_reconciliation.py | 61 ++-- .../docs/library-content-reconciliation.rst | 42 ++- .../commands/reconcile_library_search.py | 4 +- .../tests/test_content_reconciliation.py | 276 ++++++++++++++++++ .../test_content_reconciliation_adapter.py | 201 ------------- .../test_content_reconciliation_policy.py | 10 + 7 files changed, 373 insertions(+), 248 deletions(-) create mode 100644 openedx/core/djangoapps/content/search/tests/test_content_reconciliation.py delete mode 100644 openedx/core/djangoapps/content/search/tests/test_content_reconciliation_adapter.py diff --git a/openedx/core/djangoapps/content/search/api.py b/openedx/core/djangoapps/content/search/api.py index e4ef89c68564..750c6b035fd3 100644 --- a/openedx/core/djangoapps/content/search/api.py +++ b/openedx/core/djangoapps/content/search/api.py @@ -4,6 +4,7 @@ from __future__ import annotations +import json import logging import time from collections.abc import Iterator @@ -16,6 +17,7 @@ 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 @@ -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, @@ -1319,20 +1326,9 @@ def reconcile_library_components( Content publication must be quiesced for strict repair correctness; source and index rereads reduce races but cannot provide cross-system atomicity. """ - import json # pylint: disable=import-outside-toplevel - - from django.core.exceptions import ObjectDoesNotExist # pylint: disable=import-outside-toplevel - - from .content_reconciliation import reconcile_components # pylint: disable=import-outside-toplevel - from .models import SearchAccess # pylint: disable=import-outside-toplevel - if not isinstance(library_key, LibraryLocatorV2): raise ValueError("A single Libraries V2 key is required") - if ( - type(batch_size) is not int or not 1 <= batch_size <= 1000 - or type(max_documents) is not int or max_documents < 1 - ): - raise ValueError("Invalid batch_size or max_documents") + 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) @@ -1390,7 +1386,10 @@ def indexed_documents(): # truncation. JSON quoting prevents filter interpolation. while offset <= max_documents: response = index.get_documents({ - "filter": f"context_key = {json.dumps(str(library_key))} AND type = {json.dumps(DocType.library_block)}", + "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), }) diff --git a/openedx/core/djangoapps/content/search/content_reconciliation.py b/openedx/core/djangoapps/content/search/content_reconciliation.py index 069a94a86378..75c235e0cb8d 100644 --- a/openedx/core/djangoapps/content/search/content_reconciliation.py +++ b/openedx/core/djangoapps/content/search/content_reconciliation.py @@ -28,6 +28,39 @@ def as_dict(self): 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, @@ -40,20 +73,10 @@ def reconcile_components( lazy bounded streams. Full-document equality detects removed fields as drift. Index-only documents are reported, never deleted. Each scan has its own cap. """ - if type(batch_size) is not int or not 1 <= batch_size <= 1000: - raise ValueError("batch_size must be between 1 and 1000") - if type(max_documents) is not int or max_documents < 1: - raise ValueError("max_documents must be positive") + validate_reconciliation_limits(batch_size, max_documents) report = ContentDriftReport() pending = [] - def valid(doc): - 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 flush(): docs = [] for key, expected, observed in pending: @@ -75,13 +98,13 @@ def flush(): break report.scanned += 1 expected = read_source(key) - if not valid(expected) or expected["usage_key"] != str(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(observed) or observed["usage_key"] != expected["usage_key"] + not _valid_component_identity(observed, context_key) or observed["usage_key"] != expected["usage_key"] ): report.unknown += 1 continue @@ -99,12 +122,8 @@ def flush(): if pending: flush() - for position, doc in enumerate(islice(indexed_documents, max_documents + 1)): - if position == max_documents: - report.index_truncated = True - break - if not valid(doc): - report.unknown += 1 - elif read_source(doc["usage_key"]) is None: - report.index_only += 1 + _inspect_indexed_components( + report, context_key=context_key, indexed_documents=indexed_documents, + read_source=read_source, max_documents=max_documents, + ) return report diff --git a/openedx/core/djangoapps/content/search/docs/library-content-reconciliation.rst b/openedx/core/djangoapps/content/search/docs/library-content-reconciliation.rst index db343062925c..771d9d7af788 100644 --- a/openedx/core/djangoapps/content/search/docs/library-content-reconciliation.rst +++ b/openedx/core/djangoapps/content/search/docs/library-content-reconciliation.rst @@ -82,14 +82,34 @@ 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. -The isolated API adapter suite executes the actual function body with mocked -platform services and the pinned Meilisearch 0.43.0 Document/exception classes. -It checks filter scoping, page offsets/limits, full replacement, permission fields, -missing-document versus missing-index errors, rebuild checks and access metadata. -These are not full Django/database or live-engine integration tests. - -Before proposing an upstream PR, add Django library fixture tests and run the -existing search suite with a real Meilisearch. Reproduce a missed indexing event, -verify repaired content through Studio, and confirm publication quiescence -procedures with maintainers. Start with this bounded operator-only scope; propose -container/collection repair and revision-fenced background operation separately. +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. diff --git a/openedx/core/djangoapps/content/search/management/commands/reconcile_library_search.py b/openedx/core/djangoapps/content/search/management/commands/reconcile_library_search.py index 2802a3c3a42e..d43d7cd75240 100644 --- a/openedx/core/djangoapps/content/search/management/commands/reconcile_library_search.py +++ b/openedx/core/djangoapps/content/search/management/commands/reconcile_library_search.py @@ -6,6 +6,8 @@ from opaque_keys import InvalidKeyError from opaque_keys.edx.locator import LibraryLocatorV2 +from openedx.core.djangoapps.content_libraries.api import ContentLibraryNotFound + from ... import api @@ -29,6 +31,6 @@ def handle(self, *args, **options): key, repair=options["repair"], batch_size=options["batch_size"], max_documents=options["max_documents"], ) - except (InvalidKeyError, ValueError, RuntimeError) as err: + except (InvalidKeyError, ContentLibraryNotFound, ValueError, RuntimeError) as err: raise CommandError(str(err)) from err self.stdout.write(json.dumps({"repair": options["repair"], **report.as_dict()}, sort_keys=True)) diff --git a/openedx/core/djangoapps/content/search/tests/test_content_reconciliation.py b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation.py new file mode 100644 index 000000000000..ef2b13ec8447 --- /dev/null +++ b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation.py @@ -0,0 +1,276 @@ +"""CMS source/command integration and an opt-in disposable Meilisearch check.""" +# pylint: disable=protected-access +import copy +import json +import os +from io import StringIO +from unittest.mock import MagicMock, patch +from uuid import uuid4 + +import pytest +from django.conf import settings +from django.core.management import CommandError, call_command +from django.test import TestCase, override_settings +from meilisearch import Client +from meilisearch.errors import MeilisearchApiError +from meilisearch.models.document import Document, DocumentsResults +from opaque_keys.edx.locator import LibraryUsageLocatorV2 +from openedx_content import models_api as content_models +from organizations.models import Organization +from requests import Response + +from openedx.core.djangoapps.content_libraries import api as library_api +from openedx.core.djangoapps.content_tagging import api as tagging_api +from openedx.core.djangolib.testing.utils import skip_unless_cms + +try: + from .. import api + from ..documents import meili_id_from_opaque_key + from ..models import SearchAccess +except RuntimeError: + # The app is installed in CMS only; the classes below are skipped in LMS. + pass + + +def engine_error(code): + """Build the SDK exception received at the external engine boundary.""" + response = Response() + response.status_code = 404 + response._content = json.dumps({"message": code, "code": code, "type": "invalid_request"}).encode() + return MeilisearchApiError(code, response) + + +class LibraryFixtures: + # Django TestCase subclasses invoke this fixture initializer in setUp. + # pylint: disable=attribute-defined-outside-init + """Create actual library/component models; no source API or ORM doubles.""" + + def create_library_fixtures(self): + """Create canonical source records in two independent Libraries V2 contexts.""" + with override_settings(MEILISEARCH_ENABLED=False): + org = Organization.objects.create(name="Reconciliation test", short_name="repairtest") + self.library = library_api.create_library(org=org, slug="selected", title="Selected library") + self.other_library = library_api.create_library(org=org, slug="other", title="Other library") + self.first = library_api.create_library_block(self.library.key, "html", "first") + self.second = library_api.create_library_block(self.library.key, "html", "second") + self.other = library_api.create_library_block(self.other_library.key, "html", "other") + library_api.create_library_collection( + self.library.key, collection_key="collection", title="A collection", created_by=None, + ) + library_api.update_library_collection_items( + self.library.key, collection_key="collection", opaque_keys=[self.first.usage_key], + ) + taxonomy = tagging_api.create_taxonomy(name="Subject", orgs=[org], allow_multiple=True) + tagging_api.add_tag_to_taxonomy(taxonomy, tag="Algebra") + tagging_api.tag_object(str(self.first.usage_key), taxonomy, tags=["Algebra"]) + for library in (self.library, self.other_library): + SearchAccess.objects.get_or_create(context_key=library.key) + self.expected = [self.canonical(block) for block in (self.first, self.second, self.other)] + orphan_key = LibraryUsageLocatorV2(self.library.key, "html", "orphan") + self.orphan = { + **copy.deepcopy(self.expected[0]), + "id": meili_id_from_opaque_key(orphan_key), "usage_key": str(orphan_key), + } + + def canonical(self, block): + """Use the same public document serializers as existing indexing paths.""" + doc = api.searchable_doc_for_library_block(block) + doc.update(api.searchable_doc_tags(block.usage_key)) + doc.update(api.searchable_doc_collections(block.usage_key)) + doc.update(api.searchable_doc_containers(block.usage_key, "units")) + return doc + + def run_command(self, *arguments): + output = StringIO() + call_command("reconcile_library_search", str(self.library.key), *arguments, stdout=output) + return json.loads(output.getvalue()) + + +@skip_unless_cms +@override_settings(MEILISEARCH_ENABLED=True) +class TestContentReconciliation(LibraryFixtures, TestCase): + """Real CMS queries/serialization/commands with only the engine replaced.""" + + def setUp(self): + super().setUp() + self.create_library_fixtures() + self.documents = {doc["id"]: copy.deepcopy(doc) for doc in self.expected} + self.index = MagicMock(primary_key="id") + self.index.get_document.side_effect = self.read_document + self.index.get_documents.side_effect = self.read_page + self.index.add_documents.side_effect = self.replace_documents + self.client = MagicMock() + self.client.get_index.return_value = self.index + client_patch = patch.object(api, "_get_meilisearch_client", return_value=self.client) + wait_patch = patch.object(api, "_wait_for_meili_task") + client_patch.start() + wait_patch.start() + self.addCleanup(client_patch.stop) + self.addCleanup(wait_patch.stop) + self.addCleanup(content_models.Container.reset_cache) + + def read_document(self, document_id): + if document_id not in self.documents: + raise engine_error("document_not_found") + return Document(copy.deepcopy(self.documents[document_id])) + + def read_page(self, options): + docs = [doc for doc in self.documents.values() + if doc["context_key"] == str(self.library.key) and doc["type"] == "library_block"] + offset, limit = options["offset"], options["limit"] + return DocumentsResults({ + "results": copy.deepcopy(docs[offset:offset + limit]), + "offset": offset, "limit": limit, "total": len(docs), + }) + + def replace_documents(self, documents): + for doc in documents: + self.documents[doc["id"]] = copy.deepcopy(doc) + + def corrupt_documents(self): + del self.documents[self.expected[0]["id"]] + self.documents[self.expected[1]["id"]]["display_name"] = "Stale title" + self.documents[self.expected[1]["id"]]["obsolete"] = True + self.documents[self.orphan["id"]] = copy.deepcopy(self.orphan) + + def test_dry_run_reports_real_component_drift_without_writes(self): + self.corrupt_documents() + before = copy.deepcopy(self.documents) + access_count = SearchAccess.objects.count() + report = self.run_command() + assert (report["missing"], report["stale"], report["index_only"]) == (1, 1, 1) + assert report["repaired"] == 0 + assert self.documents == before + assert SearchAccess.objects.count() == access_count + + def test_repair_restores_canonical_documents_and_preserves_other_records(self): + self.corrupt_documents() + report = self.run_command("--repair", "--batch-size", "1") + assert report["repaired"] == 2 + for doc in self.expected: + assert self.documents[doc["id"]] == doc + assert self.documents[self.orphan["id"]] == self.orphan + assert self.expected[0]["tags"] + assert self.expected[0]["collections"] + assert self.expected[0]["access_id"] == SearchAccess.objects.get(context_key=self.library.key).id + clean = self.run_command() + assert (clean["missing"], clean["stale"], clean["index_only"]) == (0, 0, 1) + + def test_command_rejects_invalid_library_key(self): + with pytest.raises(CommandError): + call_command("reconcile_library_search", "not-a-library") + + def test_command_reports_missing_library_without_a_traceback(self): + with pytest.raises(CommandError): + call_command("reconcile_library_search", "lib:repairtest:missing") + + def test_caps_report_both_source_and_index_truncation(self): + report = self.run_command("--max-documents", "1", "--batch-size", "1") + assert report["scanned"] == 1 + assert report["source_truncated"] is True + assert report["index_truncated"] is True + + def test_command_refuses_actual_rebuild_lock(self): + before = copy.deepcopy(self.documents) + with api._index_rebuild_lock(api.STUDIO_LIBRARY_INDEX_NAME): + with pytest.raises(CommandError, match="rebuild in progress"): + self.run_command("--repair") + assert self.documents == before + + def test_engine_errors_abort_without_being_treated_as_missing(self): + before = copy.deepcopy(self.documents) + self.index.get_document.side_effect = engine_error("index_not_found") + with pytest.raises(MeilisearchApiError): + self.run_command("--repair") + assert self.documents == before + + def test_missing_access_metadata_is_not_created_by_dry_run(self): + SearchAccess.objects.filter(context_key=self.library.key).delete() + with pytest.raises(CommandError, match="access metadata is missing"): + self.run_command() + assert not SearchAccess.objects.filter(context_key=self.library.key).exists() + + def test_missing_index_aborts_instead_of_reporting_missing_content(self): + self.client.get_index.side_effect = engine_error("index_not_found") + with pytest.raises(MeilisearchApiError): + self.run_command() + + def test_wrong_primary_key_refuses_repair(self): + before = copy.deepcopy(self.documents) + self.index.primary_key = "usage_key" + with pytest.raises(CommandError, match="Reconcile index settings"): + self.run_command("--repair") + assert self.documents == before + + def test_corrupt_indexed_usage_key_is_unknown_and_preserved(self): + malformed = {**self.orphan, "usage_key": "invalid-key"} + self.documents[malformed["id"]] = malformed + report = self.run_command("--repair") + assert report["unknown"] == 1 + assert report["index_only"] == 0 + assert self.documents[malformed["id"]] == malformed + + def test_primary_key_collision_cannot_replace_another_library_document(self): + collision = {**self.expected[2], "id": self.expected[0]["id"]} + self.documents[collision["id"]] = collision + report = self.run_command("--repair") + assert report["unknown"] == 1 + assert report["repaired"] == 0 + assert self.documents[collision["id"]] == collision + + def test_duplicate_usage_key_with_wrong_id_is_unknown_and_preserved(self): + malformed = {**self.expected[0], "id": "noncanonical-id"} + self.documents[malformed["id"]] = malformed + report = self.run_command("--repair") + assert report["unknown"] == 1 + assert report["repaired"] == 0 + assert self.documents[malformed["id"]] == malformed + self.index.add_documents.assert_not_called() + + def test_failed_engine_write_aborts_repair(self): + self.corrupt_documents() + before = copy.deepcopy(self.documents) + self.index.add_documents.side_effect = engine_error("internal") + with pytest.raises(MeilisearchApiError): + self.run_command("--repair") + assert self.documents == before + + +@skip_unless_cms +@pytest.mark.skipif(not os.environ.get("OPENEDX_SEARCH_LIVE_URL"), reason="Disposable Meilisearch URL not supplied") +@override_settings(MEILISEARCH_ENABLED=True) +class TestLiveContentReconciliation(LibraryFixtures, TestCase): + """One small, explicit live-engine repair with no background mutations.""" + + def test_live_missing_stale_orphan_and_other_library(self): + self.create_library_fixtures() + self.addCleanup(content_models.Container.reset_cache) + client = Client( + os.environ["OPENEDX_SEARCH_LIVE_URL"], + os.environ.get("OPENEDX_SEARCH_LIVE_KEY", settings.MEILISEARCH_API_KEY), + ) + name = "reconciliation_test_" + uuid4().hex + client.wait_for_task(client.create_index(name, {"primaryKey": "id"}).task_uid) + try: + index = client.get_index(name) + client.wait_for_task(index.update_filterable_attributes(["context_key", "type"]).task_uid) + client.wait_for_task(index.add_documents(self.expected).task_uid) + client.wait_for_task(index.delete_document(self.expected[0]["id"]).task_uid) + stale = {**self.expected[1], "display_name": "Stale title", "obsolete": True} + client.wait_for_task(index.add_documents([stale, self.orphan]).task_uid) + before = [vars(index.get_document(doc["id"])) for doc in (stale, self.orphan, self.expected[2])] + with patch.object(api, "_get_meilisearch_client", return_value=client), \ + patch.object(api, "STUDIO_LIBRARY_INDEX_NAME", name): + dry = self.run_command() + assert (dry["missing"], dry["stale"], dry["index_only"]) == (1, 1, 1) + assert dry["repaired"] == 0 + assert [vars(index.get_document(doc["id"])) + for doc in (stale, self.orphan, self.expected[2])] == before + repaired = self.run_command("--repair", "--batch-size", "1") + assert repaired["repaired"] == 2 + clean = self.run_command() + assert (clean["missing"], clean["stale"], clean["index_only"]) == (0, 0, 1) + for doc in (*self.expected, self.orphan): + assert vars(index.get_document(doc["id"])) == doc + finally: + client.wait_for_task(client.delete_index(name).task_uid) diff --git a/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_adapter.py b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_adapter.py deleted file mode 100644 index d0c6f28fba69..000000000000 --- a/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_adapter.py +++ /dev/null @@ -1,201 +0,0 @@ -"""Isolated adapter contract tests with the pinned Meilisearch SDK. - -Execute the real API function body with mocked platform services. This validates -adapter behavior, not Django ORM/database integration or a live Meilisearch. -""" -# ruff: noqa: PT009, PT027 -import ast -import importlib.util -import sys -import types -import unittest -from pathlib import Path -from unittest.mock import MagicMock, patch - -from meilisearch.errors import MeilisearchApiError -from meilisearch.models.document import Document, DocumentsResults - -ROOT = Path(__file__).parents[1] -SPEC = importlib.util.spec_from_file_location("isolated_search.content_reconciliation", ROOT / "content_reconciliation.py") -POLICY = importlib.util.module_from_spec(SPEC) -sys.modules[SPEC.name] = POLICY -SPEC.loader.exec_module(POLICY) - - -class LibraryKey: - """Minimal stand-in for an already-parsed library key.""" - - def __str__(self): - return "lib:org:one" - - -class ComponentKey: - """Minimal usage key satisfying the API's parsing/context contract.""" - - def __init__(self, name): - self.name = name - self.context_key = LIBRARY - - def __str__(self): - return self.name - - @classmethod - def from_string(cls, key): - return cls(key) - - -LIBRARY = LibraryKey() - - -def doc(key="a", **fields): - return {"id": key, "usage_key": key, "context_key": str(LIBRARY), "type": "library_block", **fields} - - -class NotFound(Exception): - """Stand-in for Django ObjectDoesNotExist.""" - - -def engine_error(code): - # Use the real SDK exception class, without constructing an HTTP response. - err = MeilisearchApiError.__new__(MeilisearchApiError) - Exception.__init__(err, code) - err.code = code - return err - - -class AdapterTest(unittest.TestCase): - """Verify the actual API adapter, with only external dependencies replaced.""" - - def setUp(self): - self.lib = MagicMock() - self.index = MagicMock(primary_key="id") - self.client = MagicMock() - self.client.get_index.return_value = self.index - self.access = MagicMock() - self.access.objects.filter.return_value.exists.return_value = True - self.lib.get_library_components.return_value.order_by.return_value.iterator.return_value = ["a"] - self.lib.get_library_components.return_value.filter.return_value.exists.return_value = True - self.lib.get_component_from_usage_key.side_effect = lambda key: types.SimpleNamespace(pk=str(key)) - self.lib.LibraryXBlockMetadata.from_component.side_effect = lambda key, comp: types.SimpleNamespace( - usage_key=ComponentKey(comp if isinstance(comp, str) else comp.pk), - ) - self.lib.ContentLibraryBlockNotFound = NotFound - self.index.get_document.return_value = Document(doc()) - self.index.get_documents.return_value = DocumentsResults({"results": [], "offset": 0, "limit": 100, "total": 0}) - self.rebuild = MagicMock(return_value=None) - self.wait = MagicMock() - self.namespace = { - "__package__": "isolated_search", "LibraryLocatorV2": LibraryKey, "UsageKey": ComponentKey, - "lib_api": self.lib, "STUDIO_LIBRARY_INDEX_NAME": "library", "INDEX_PRIMARY_KEY": "id", - "DocType": types.SimpleNamespace(library_block="library_block"), - "Fields": types.SimpleNamespace(usage_key="usage_key"), "InvalidKeyError": ValueError, - "MeilisearchApiError": MeilisearchApiError, "_get_meilisearch_client": lambda: self.client, - "_get_running_rebuild_index_name": self.rebuild, "_wait_for_meili_task": self.wait, - "searchable_doc_for_library_block": lambda metadata: doc(str(metadata.usage_key), access_id=7), - "searchable_doc_tags": lambda key: {"tags": {}}, - "searchable_doc_collections": lambda key: {"collections": {}}, - "searchable_doc_containers": lambda key, group: {group: {}}, - } - tree = ast.parse((ROOT / "api.py").read_text()) - target = next(node for node in tree.body if isinstance(node, ast.FunctionDef) - and node.name == "reconcile_library_components") - exec(compile(ast.Module(body=[target], type_ignores=[]), "api.py", "exec"), self.namespace) - self.function = self.namespace["reconcile_library_components"] - exceptions = types.ModuleType("django.core.exceptions") - exceptions.ObjectDoesNotExist = NotFound - models = types.ModuleType("isolated_search.models") - models.SearchAccess = self.access - self.modules = patch.dict(sys.modules, {"django.core.exceptions": exceptions, "isolated_search.models": models}) - self.modules.start() - self.addCleanup(self.modules.stop) - - def test_dry_run_document_sdk_and_no_writes(self): - report = self.function(LIBRARY) - self.assertEqual(report.stale, 1) - self.index.add_documents.assert_not_called() - self.access.objects.create.assert_not_called() - self.assertEqual(self.index.get_documents.call_args.args[0]["filter"], - 'context_key = "lib:org:one" AND type = "library_block"') - - def test_repair_full_document_and_wait(self): - report = self.function(LIBRARY, repair=True) - self.assertEqual(report.repaired, 1) - written = self.index.add_documents.call_args.args[0][0] - self.assertEqual(written["access_id"], 7) - self.assertEqual(written["tags"], {}) - self.assertEqual(written["units"], {}) - self.index.update_documents.assert_not_called() - self.wait.assert_called_once_with(self.index.add_documents.return_value) - - def test_only_document_not_found_is_absence(self): - self.index.get_document.side_effect = engine_error("document_not_found") - self.assertEqual(self.function(LIBRARY).missing, 1) - - def test_missing_index_not_treated_as_missing_document(self): - self.index.get_document.side_effect = engine_error("index_not_found") - with self.assertRaises(MeilisearchApiError): - self.function(LIBRARY, repair=True) - self.index.add_documents.assert_not_called() - - def test_engine_authorization_failure_aborts(self): - self.index.get_document.side_effect = engine_error("invalid_api_key") - with self.assertRaises(MeilisearchApiError): - self.function(LIBRARY, repair=True) - self.index.add_documents.assert_not_called() - - def test_active_rebuild_refused(self): - self.rebuild.return_value = "library_new" - with self.assertRaises(RuntimeError): - self.function(LIBRARY, repair=True) - self.index.add_documents.assert_not_called() - - def test_rebuild_started_before_write_refused(self): - self.rebuild.side_effect = [None, "library_new"] - with self.assertRaises(RuntimeError): - self.function(LIBRARY, repair=True) - self.index.add_documents.assert_not_called() - - def test_missing_access_dry_run_refused_without_creation(self): - self.access.objects.filter.return_value.exists.return_value = False - with self.assertRaises(ValueError): - self.function(LIBRARY) - self.client.get_index.assert_not_called() - - def test_paginated_documents_bounded_offsets(self): - self.lib.get_library_components.return_value.order_by.return_value.iterator.return_value = [] - self.index.get_documents.side_effect = [ - DocumentsResults({"results": [doc("a"), doc("b")], "offset": 0, "limit": 2, "total": 4}), - DocumentsResults({"results": [doc("c")], "offset": 2, "limit": 1, "total": 4}), - ] - report = self.function(LIBRARY, batch_size=2, max_documents=2) - self.assertTrue(report.index_truncated) - self.assertEqual([call.args[0]["offset"] for call in self.index.get_documents.call_args_list], [0, 2]) - self.assertEqual([call.args[0]["limit"] for call in self.index.get_documents.call_args_list], [2, 1]) - - def test_malformed_index_key_reported_unknown(self): - self.lib.get_library_components.return_value.order_by.return_value.iterator.return_value = [] - self.index.get_documents.side_effect = [ - DocumentsResults({"results": [doc(usage_key=42)], "offset": 0, "limit": 100, "total": 1}), - DocumentsResults({"results": [], "offset": 1, "limit": 100, "total": 1}), - ] - report = self.function(LIBRARY) - self.assertEqual(report.unknown, 1) - self.assertEqual(report.index_only, 0) - - def test_wrong_primary_key_refused(self): - self.index.primary_key = "other" - with self.assertRaises(ValueError): - self.function(LIBRARY, repair=True) - self.index.add_documents.assert_not_called() - - def test_builder_failure_propagates_not_orphan(self): - def fail_builder(_): - raise NotFound("related source read failed") - self.namespace["searchable_doc_for_library_block"] = fail_builder - with self.assertRaises(NotFound): - self.function(LIBRARY, repair=True) - self.index.add_documents.assert_not_called() - - -if __name__ == "__main__": - unittest.main() diff --git a/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_policy.py b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_policy.py index 463e4db39004..8b5c7b83d319 100644 --- a/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_policy.py +++ b/openedx/core/djangoapps/content/search/tests/test_content_reconciliation_policy.py @@ -24,7 +24,11 @@ def document(key="a", **fields): class ReconcilePolicyTest(unittest.TestCase): """Exercise bounded writes, ambiguous identity and observed mutation races.""" + def setUp(self): + self.writes = [] + def run_reconcile(self, sources=None, indexed=None, **kwargs): + """Run the policy against in-memory source/index boundaries.""" sources = {"a": document()} if sources is None else sources indexed = {} if indexed is None else indexed self.writes = [] @@ -41,6 +45,12 @@ def write_batch(docs): options.update(kwargs) return reconcile_components(**options) + def test_indexed_usage_key_with_noncanonical_id_is_unknown(self): + report = self.run_reconcile(indexed={"noncanonical": document(id="noncanonical")}, repair=True) + self.assertEqual(report.unknown, 1) + self.assertEqual(report.repaired, 1) # Restore the absent canonical ID only. + self.assertEqual(self.writes, [[document()]]) + def test_dry_run_default(self): report = self.run_reconcile() self.assertEqual(report.missing, 1)