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
50 changes: 50 additions & 0 deletions docs/contributions/async-library-indexing.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
# Durable library-wide indexing slice

## Baseline and boundary

Base: `7944dc908dffff5ff21754859fc0ca554005b313` of openedx/edx-platform.
Content search already has Celery tasks. Library create and rename handlers call their library-wide task with `.apply`, which executes synchronously. Some source library events are themselves emitted by workers; this change decouples that event worker from the engine wait, rather than asserting every event originates in an HTTP publication request.

## Implemented contract

`LIBRARY_SEARCH_ASYNC_INDEXING=True` opts only library create/rename into a durable database intent. Default behavior stays compatible with the current frontend. SearchAccess creation remains immediate. Library block, container, collection, association and course flows retain their existing behavior.

Within the caller's transaction, request_library_index increments a per-library requested revision. A post-commit callback dispatches a Celery task. Broker operational failures are logged; the database request survives. A rolled-back caller transaction leaves neither intent nor dispatch.

The worker takes the library request row lock, reads current authoritative content rather than an old event payload, streams block, collection and container documents through the existing document constructors and `_update_index_docs` API boundary (including their breadcrumbs), then advances completion only after every batch's engine task wait returns. Source querysets use `iterator(chunk_size=batch_size)` rather than populating Django's full result cache. The default batch size is 100 documents, configurable with the positive integer `LIBRARY_SEARCH_INDEX_BATCH_SIZE`. The synchronous baseline API remains unchanged. Each batch retains the existing current-index wait and concurrent-rebuild dual-write behavior; the existing rebuild shadow write is not newly given a completion wait. Duplicate deliveries are no-ops; multiple pending requests coalesce. Failures roll back completion. Celery retries connection/Meilisearch errors five times with backoff and jitter, and late acknowledgment plus worker-loss rejection allow redelivery. A bounded recovery command redispatches unresolved rows, including requests lost between database commit and broker dispatch. Run it periodically through the deployment's existing scheduler; no scheduler is silently installed.

Library deletion cancels the durable row. A missing-library exception during delivery also removes the request, avoiding a permanent poison item. This cancellation does not introduce deletion of existing search documents; library lifecycle cleanup remains the existing upstream responsibility.

## Lock and ordering limits

The row lock spans the engine call deliberately: queued library-wide writes cannot overtake one another, and enqueue revisions cannot be acknowledged accidentally by an older worker. A crash after an engine write but before the database commit causes replay of current content, giving idempotent convergence rather than exactly-once delivery.

This conservative design means a subsequent create/rename enqueue for the same library can wait for an active worker's engine call. It moves the ordinary wait out of the event path but does not guarantee bounded authoring latency under lock contention. Pending revisions are ordered by transaction arrival; they are not source content versions. They always reread current content, so delayed notifications converge on current state. Existing per-item writes do not acquire this lock: their concurrent changes can race library-wide indexing. Therefore global stale-write prevention is not claimed.

SQLite does not enforce select_for_update. The isolated tests prove persistence/rollback/coalescing, not production row-lock ordering. Before activation, run concurrent-worker/rename/deletion fault tests against the supported production database and the real engine. For nonblocking enqueue, a future immutable-intent table with a separate per-library worker lock is preferable; it needs explicit FK-lock and deletion semantics on the deployment database.

## Rollout and outstanding work

1. Obtain maintainer agreement for the bounded scope and setting name; run CMS integration tests and migration checks with pinned dependencies.
2. Deploy the additive migration with the setting false. Deploy workers before enabling the setting.
3. Configure periodic bounded recovery and alert on age/revision lag; validate engine timeout and broker retry settings.
4. Add an author-visible pending/failed indicator and refetch behavior before opt-in production use. This slice exposes durable internal state, not a new public status endpoint or frontend.
5. Run production-database concurrency tests and a real Meilisearch fault-injection scenario. Redesign the long lock or unify per-item writes before claiming comprehensive order safety.

Document batches are bounded, but document byte size, nested metadata construction, and database-driver buffering are not. End-to-end process memory therefore requires measurement on the production database. Engine failure after an accepted batch leaves the revision pending; recovery rebuilds batches from current content and replays accepted earlier batches. Library disappearance after a partial write cancels the request, but does not clean up partially indexed documents. Library deletion/reconciliation must handle that existing lifecycle gap.

Per-item durable intents, source versions, cross-engine support, telemetry, job dead-letter UI, library-wide deletion repair, retention cleanup, and deployment benchmarks remain open. This is an initial contribution for library-wide indexing requests.

Disabling the setting restores synchronous future events but does not cancel already pending requests. Drain/recover pending rows before rollback; workers can continue draining with the setting false. Do not drop the table with pending work.

## Local validation evidence

17 focused tests passed with complete CMS dependencies. Django reported no missing
search migrations. The actual additive migration applied and rolled back on a
dedicated MySQL 8.4.11 database using the existing migration dependency state.
That same database preserved four concurrent same-library enqueues, serialized
duplicate deliveries into one indexing call, and recovered a pending revision
after its task process was killed. A real Redis broker connection outage left
the intent pending; a Celery worker consumed it after recovery and ignored a
duplicate. Those fault scenarios mock the indexing/engine boundary. They prove
outbox and broker behavior, not complete source-to-engine concurrency safety.
15 changes: 12 additions & 3 deletions openedx/core/djangoapps/content/search/handlers.py
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

import logging

from django.conf import settings
from django.db.models.signals import post_delete
from django.dispatch import receiver
from opaque_keys import InvalidKeyError
Expand Down Expand Up @@ -42,7 +43,7 @@
)

from openedx.core.djangoapps.content.course_overviews.models import CourseOverview
from openedx.core.djangoapps.content.search.models import SearchAccess
from openedx.core.djangoapps.content.search.models import LibraryIndexRequest, SearchAccess
from openedx.core.djangoapps.content_libraries import api as lib_api
from xmodule.modulestore.django import SignalHandler

Expand All @@ -54,6 +55,7 @@
upsert_item_collections_index_docs,
upsert_item_containers_index_docs,
)
from .library_indexing import request_library_index
from .tasks import (
delete_course_index_docs,
delete_library_block_index_doc,
Expand Down Expand Up @@ -111,6 +113,7 @@ def delete_course_search_access(sender, instance, **kwargs): # pylint: disable=
@receiver(CONTENT_LIBRARY_DELETED)
def delete_library_search_access(content_library: ContentLibraryData, **kwargs):
"""Deletes the SearchAccess instance for deleted content libraries"""
LibraryIndexRequest.objects.filter(library_key=str(content_library.library_key)).delete()
SearchAccess.objects.filter(context_key=content_library.library_key).delete()


Expand Down Expand Up @@ -237,7 +240,10 @@ def content_library_created_handler(**kwargs) -> None:
# right after creation. Without this, the JWT token won't include the new library's
# access_id until it's added by the document indexing process or the page is refreshed.
SearchAccess.objects.get_or_create(context_key=library_key)
update_content_library_index_docs.apply(args=[str(library_key), True])
if getattr(settings, "LIBRARY_SEARCH_ASYNC_INDEXING", False):
request_library_index(library_key)
else:
update_content_library_index_docs.apply(args=[str(library_key), True])


@receiver(CONTENT_LIBRARY_UPDATED)
Expand All @@ -257,7 +263,10 @@ def content_library_updated_handler(**kwargs) -> None:
# Update ALL items in the library, because their breadcrumbs will be outdated.
# TODO: just patch the "breadcrumbs" field? It's the same on every one.
# TODO: check if the library display_name has actually changed before updating all items?
update_content_library_index_docs.apply(args=[str(library_key)])
if getattr(settings, "LIBRARY_SEARCH_ASYNC_INDEXING", False):
request_library_index(library_key)
else:
update_content_library_index_docs.apply(args=[str(library_key)])


@receiver(LIBRARY_COLLECTION_CREATED)
Expand Down
71 changes: 71 additions & 0 deletions openedx/core/djangoapps/content/search/library_indexing.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,71 @@
"""Durable library-wide indexing with an explicit, recoverable broker boundary."""

import logging

from celery import shared_task
from django.conf import settings
from django.db import transaction
from kombu.exceptions import OperationalError
from meilisearch.errors import MeilisearchError
from opaque_keys.edx.locator import LibraryLocatorV2

from .library_indexing_batches import index_library_in_batches
from .models import LibraryIndexRequest

log = logging.getLogger(__name__)


def _dispatch(library_key):
"""A broker outage must not undo a committed content change."""
try:
process_library_index_request.delay(library_key)
except OperationalError:
log.exception("Library indexing dispatch failed; durable request remains pending for %s", library_key)


def request_library_index(library_key):
"""Record intent in the caller's transaction and dispatch only after commit."""
library_key = str(LibraryLocatorV2.from_string(str(library_key)))
with transaction.atomic():
request, _ = LibraryIndexRequest.objects.get_or_create(library_key=library_key)
request = LibraryIndexRequest.objects.select_for_update().get(pk=request.pk)
request.requested_revision += 1
request.save(update_fields=["requested_revision", "updated_at"])
transaction.on_commit(lambda: _dispatch(library_key))


@shared_task(
autoretry_for=(MeilisearchError, ConnectionError),
retry_backoff=True,
retry_jitter=True,
max_retries=5,
acks_late=True,
reject_on_worker_lost=True,
)
def process_library_index_request(library_key):
"""Serialize writes for one library and acknowledge only completed engine tasks.

The conservative row lock spans the engine call. Duplicate broker deliveries
become no-ops; a crashed worker leaves the revision pending. This lock requires
a database supporting row locks (the production MySQL/PostgreSQL backends).
"""
from openedx.core.djangoapps.content_libraries.api import ( # pylint: disable=import-outside-toplevel
ContentLibraryNotFound,
)


with transaction.atomic():
request = LibraryIndexRequest.objects.select_for_update().filter(library_key=library_key).first()
if request is None or request.completed_revision == request.requested_revision:
return
try:
index_library_in_batches(
LibraryLocatorV2.from_string(library_key),
batch_size=getattr(settings, "LIBRARY_SEARCH_INDEX_BATCH_SIZE", 100),
)
except ContentLibraryNotFound:
# Deletion may precede delivery even if its cancellation signal was lost.
request.delete()
return
request.completed_revision = request.requested_revision
request.save(update_fields=["completed_revision", "updated_at"])
36 changes: 36 additions & 0 deletions openedx/core/djangoapps/content/search/library_indexing_batches.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
"""Stream full library documents into bounded, individually awaited engine batches."""

from openedx.core.djangoapps.content_libraries import api as lib_api


def index_library_in_batches(library_key, batch_size=100):
"""Reuse existing document construction and dual-index write behavior.

Each source queryset uses iterator() to avoid Django's full result cache.
No persistent per-batch checkpoint is used: interrupted runs replay current
authoritative content, including batches the engine may already have accepted.
"""
if not isinstance(batch_size, int) or isinstance(batch_size, bool) or batch_size <= 0:
raise ValueError("batch_size must be a positive integer")

from . import api # pylint: disable=import-outside-toplevel

def documents():
for component in lib_api.get_library_components(library_key).iterator(chunk_size=batch_size):
metadata = lib_api.LibraryXBlockMetadata.from_component(library_key, component)
yield api.searchable_doc_for_library_block(metadata)
for container in lib_api.get_library_containers(library_key).iterator(chunk_size=batch_size):
container_key = lib_api.library_container_locator(library_key, container)
yield api.searchable_doc_for_container(container_key)
for collection in lib_api.get_library_collections(library_key).iterator(chunk_size=batch_size):
collection_key = lib_api.library_collection_locator(library_key, collection.collection_code)
yield api.searchable_doc_for_collection(collection_key, collection=collection)

batch = []
for doc in documents():
batch.append(doc)
if len(batch) == batch_size:
api._update_index_docs(api.STUDIO_LIBRARY_INDEX_NAME, batch) # pylint: disable=protected-access
batch = []
if batch:
api._update_index_docs(api.STUDIO_LIBRARY_INDEX_NAME, batch) # pylint: disable=protected-access
Original file line number Diff line number Diff line change
@@ -0,0 +1,27 @@
"""Redispatch durable requests after a broker outage or exhausted worker retries."""

from django.core.management.base import BaseCommand
from django.db.models import F

from ...library_indexing import process_library_index_request
from ...models import LibraryIndexRequest


class Command(BaseCommand):
"""Dispatch only the requested bounded number of pending revisions."""
help = "Redispatch a bounded batch of pending library indexing requests."

def add_arguments(self, parser):
parser.add_argument("--limit", type=int, default=100)

def handle(self, *args, **options):
if options["limit"] <= 0:
raise ValueError("--limit must be positive")
pending = LibraryIndexRequest.objects.filter(
requested_revision__gt=F("completed_revision"),
).order_by("updated_at").values_list("library_key", flat=True)[:options["limit"]]
count = 0
for library_key in pending:
process_library_index_request.delay(library_key)
count += 1
self.stdout.write(f"Dispatched {count} pending library indexing requests.")
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
from django.db import migrations, models


class Migration(migrations.Migration):
dependencies = [("search", "0003_clear_library_incremental_index_checkpoints")]

operations = [
migrations.CreateModel(
name="LibraryIndexRequest",
fields=[
("id", models.BigAutoField(auto_created=True, primary_key=True, serialize=False, verbose_name="ID")),
("library_key", models.CharField(max_length=255, unique=True)),
("requested_revision", models.PositiveBigIntegerField(default=0)),
("completed_revision", models.PositiveBigIntegerField(default=0)),
("updated_at", models.DateTimeField(auto_now=True)),
],
),
]
12 changes: 12 additions & 0 deletions openedx/core/djangoapps/content/search/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -81,3 +81,15 @@ class IncrementalIndexCompleted(models.Model): # noqa: DJ008
unique=True,
null=False,
)


class LibraryIndexRequest(models.Model): # noqa: DJ008
"""Durable, coalesced library-wide indexing intent. Contains no learner data.

.. no_pii:
"""

library_key = models.CharField(max_length=255, unique=True)
requested_revision = models.PositiveBigIntegerField(default=0)
completed_revision = models.PositiveBigIntegerField(default=0)
updated_at = models.DateTimeField(auto_now=True)
1 change: 1 addition & 0 deletions openedx/core/djangoapps/content/search/tasks.py
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
)

from . import api
from .library_indexing import process_library_index_request # noqa: F401 # pylint: disable=unused-import

log = logging.getLogger(__name__)

Expand Down
Loading