From 38e6ee21ed5a640b36ae375c4307965579cd4e07 Mon Sep 17 00:00:00 2001 From: Christopher Patti Date: Fri, 17 Jul 2026 17:56:27 -0400 Subject: [PATCH] Serialize new resource embedding fan-out --- vector_search/tasks.py | 7 ++++--- vector_search/tasks_test.py | 17 ++++++++++++++++- 2 files changed, 20 insertions(+), 4 deletions(-) diff --git a/vector_search/tasks.py b/vector_search/tasks.py index 4e842589bc..9c3b34d9b0 100644 --- a/vector_search/tasks.py +++ b/vector_search/tasks.py @@ -433,9 +433,10 @@ def embed_new_learning_resources(self): ) ] ) - embed_tasks = celery.group(tasks) - - return self.replace(embed_tasks) + # Dispatch chunks sequentially instead of materializing the whole lookback + # window as one group. This bounds broker queue depth and prevents KEDA from + # scaling the embeddings worker fleet for a short-lived fan-out burst. + return _replace_with_chain(self, tasks) @app.task(bind=True) diff --git a/vector_search/tasks_test.py b/vector_search/tasks_test.py index 769339ba19..9419b1479b 100644 --- a/vector_search/tasks_test.py +++ b/vector_search/tasks_test.py @@ -235,13 +235,28 @@ def test_embed_new_learning_resources(mocker, mocked_celery): with pytest.raises(mocked_celery.replace_exception_class): embed_new_learning_resources.delay() - list(mocked_celery.group.call_args[0][0]) assert generate_embeddings_mock.si.call_count == 1 + embedding_signature = generate_embeddings_mock.si.return_value + mocked_celery.chain.assert_called_once_with(embedding_signature) + mocked_celery.group.assert_not_called() + assert mocked_celery.replace.call_args[0][1] == mocked_celery.chain.return_value embedded_ids = generate_embeddings_mock.si.mock_calls[0].args[0] assert sorted(new_resource_ids) == sorted(embedded_ids) +def test_embed_new_learning_resources_no_work(mocker, mocked_celery): + """embed_new_learning_resources should not dispatch an empty canvas""" + settings.QDRANT_EMBEDDINGS_TASK_LOOKBACK_WINDOW = 1 + mocker.patch("vector_search.tasks.now_in_utc", return_value=now_in_utc()) + + assert embed_new_learning_resources.run() is None + + mocked_celery.chain.assert_not_called() + mocked_celery.group.assert_not_called() + mocked_celery.replace.assert_not_called() + + def test_embed_new_content_files(mocker, mocked_celery): """ embed_new_content_files should generate embeddings for new content files