diff --git a/src/backend/core/tasks/search.py b/src/backend/core/tasks/search.py index e1c39e6b..083e2e03 100644 --- a/src/backend/core/tasks/search.py +++ b/src/backend/core/tasks/search.py @@ -20,7 +20,12 @@ logger = getLogger(__file__) @app.task def document_indexer_task(document_id): - """Celery Task : Sends indexation query for a document.""" + """ + Celery Task: Indexes a single document by its ID. + + Args: + document_id: Primary key of the document to index. + """ indexer = get_document_indexer() if indexer: @@ -30,8 +35,17 @@ def document_indexer_task(document_id): def batch_indexer_throttle_acquire(timeout: int = 0, atomic: bool = True): """ - Enable the task throttle flag for a delay. - Uses redis locks if available to ensure atomic changes + Acquire a throttle lock to prevent multiple batch indexation tasks during countdown. + + implements a debouncing pattern: only the first call during the timeout period + will succeed, subsequent calls are skipped until the timeout expires. + + Args: + timeout (int): Lock duration in seconds (countdown period). + atomic (bool): Use Redis locks for atomic operations if available. + + Returns: + bool: True if lock acquired (first call), False if already held (subsequent calls). """ key = "document-batch-indexer-throttle" @@ -41,15 +55,23 @@ def batch_indexer_throttle_acquire(timeout: int = 0, atomic: bool = True): with cache.locks(key): return batch_indexer_throttle_acquire(timeout, atomic=False) - # Use add() here : - # - set the flag and returns true if not exist - # - do nothing and return false if exist + # cache.add() is atomic test-and-set operation: + # - If key doesn't exist: creates it with timeout and returns True + # - If key already exists: does nothing and returns False + # The key expires after timeout seconds, releasing the lock. + # The value 1 is irrelevant, only the key presence/absence matters. return cache.add(key, 1, timeout=timeout) @app.task def batch_document_indexer_task(timestamp): - """Celery Task : Sends indexation query for a batch of documents.""" + """ + Celery Task: Batch indexes all documents modified since timestamp. + + Args: + timestamp: ISO timestamp to filter documents by updated_at, deleted_at, + or ancestors_deleted_at. + """ indexer = get_document_indexer() if indexer: @@ -65,20 +87,29 @@ def batch_document_indexer_task(timestamp): def trigger_batch_document_indexer(document): """ - Trigger indexation task with debounce a delay set by the SEARCH_INDEXER_COUNTDOWN setting. - + Trigger document indexation with optional debounce mechanism. + + behavior depends on SEARCH_INDEXER_COUNTDOWN setting: + - if countdown > 0 sec (async batch mode): + * schedules a batch indexation task after countdown in seconds + * uses throttle mechanism to ensure only ONE batch task runs per countdown period + * all documents modified since first trigger are indexed together + - if countdown == 0 sec (sync mode): + * executes indexation synchronously in the current thread + * no batching, no throttling, no Celery task queuing + Args: - document (Document): The document instance. + document (Document): the document instance that triggered the indexation. """ countdown = int(settings.SEARCH_INDEXER_COUNTDOWN) - # DO NOT create a task if indexation if disabled + # DO NOT create a task if indexation is disabled if not settings.SEARCH_INDEXER_CLASS: return if countdown > 0: - # Each time this method is called during a countdown, we increment the - # counter and each task decrease it, so the index be run only once. + # use throttle to ensure only one task is scheduled per countdown period. + # if throttle acquired, schedule batch task; otherwise skip. if batch_indexer_throttle_acquire(timeout=countdown): logger.info( "Add task for batch document indexation from updated_at=%s in %d seconds",