diff --git a/modules/weko-search-ui/tests/conftest.py b/modules/weko-search-ui/tests/conftest.py index dd3cc375e1..f6a33c5bae 100644 --- a/modules/weko-search-ui/tests/conftest.py +++ b/modules/weko-search-ui/tests/conftest.py @@ -707,7 +707,9 @@ def base_app(instance_path, search_class, request): "internal report": "other", "report part": "other", "conference object": "conference output", - } + }, + WEKO_SEARCH_UI_CELERY_STATUS_CACHE_TTL = 60, + WEKO_SEARCH_UI_CELERY_STATUS = "weko_search_ui_celery_status" ) app_.url_map.converters["pid"] = PIDConverter app_.config["RECORDS_REST_ENDPOINTS"]["recid"]["search_class"] = search_class diff --git a/modules/weko-search-ui/tests/test_admin.py b/modules/weko-search-ui/tests/test_admin.py index 8b371e946e..78df596158 100644 --- a/modules/weko-search-ui/tests/test_admin.py +++ b/modules/weko-search-ui/tests/test_admin.py @@ -903,7 +903,9 @@ class TestItemBulkExport: # .tox/c1/bin/pytest --cov=weko_search_ui tests/test_admin.py::TestItemBulkExport::test_check_export_status -vv -s --cov-branch --cov-report=term --basetemp=/code/modules/weko-search-ui/.tox/c1/tmp def test_check_export_status(self,app,client,users, redis_connect,mocker): - mocker.patch("weko_search_ui.admin.check_celery_is_run",return_value=True) + mock_check_celery_is_run = mocker.patch( + "weko_search_ui.admin.check_celery_is_run", return_value=True + ) mocker.patch("weko_search_ui.admin.check_session_lifetime",return_value=True) start_time_str = '2024/05/01 12:55:36' @@ -927,6 +929,7 @@ def test_check_export_status(self,app,client,users, redis_connect,mocker): 'status': 'STARTED', 'uri_status': False }} + mock_check_celery_is_run.assert_called_with(is_task=True) with patch('weko_search_ui.admin.get_export_status', return_value=(True, 'test_uri', '', '', 'STARTED', start_time_str, '')): @@ -942,6 +945,7 @@ def test_check_export_status(self,app,client,users, redis_connect,mocker): 'status': 'STARTED', 'uri_status': True }} + mock_check_celery_is_run.assert_called_with(is_task=True) # .tox/c1/bin/pytest --cov=weko_search_ui tests/test_admin.py::TestItemBulkExport::test_cancel_export -vv -s --cov-branch --cov-report=term --basetemp=/code/modules/weko-search-ui/.tox/c1/tmp def test_cancel_export(self, app, client, users, redis_connect, mocker): diff --git a/modules/weko-search-ui/tests/test_tasks.py b/modules/weko-search-ui/tests/test_tasks.py index d5cf6138c8..cd800f3467 100644 --- a/modules/weko-search-ui/tests/test_tasks.py +++ b/modules/weko-search-ui/tests/test_tasks.py @@ -6,6 +6,7 @@ import pytest import unittest from flask import current_app +from redis import RedisError from unittest.mock import patch, MagicMock from flask_login import current_user @@ -422,11 +423,41 @@ def test_is_import_running(i18n_app): # def check_celery_is_run(): # .tox/c1/bin/pytest --cov=weko_search_ui tests/test_tasks.py::test_check_celery_is_run -vv -s --cov-branch --cov-report=term --basetemp=/code/modules/weko-search-ui/.tox/c1/tmp def test_check_celery_is_run(i18n_app): - with patch("celery.task.control.inspect.ping",return_value={'hostname': True}): - assert check_celery_is_run()==True - - with patch("celery.task.control.inspect.ping",return_value={}): - assert check_celery_is_run()==False + cache_key = "weko_search_ui_celery_status" + datastore = MagicMock() + + is_task = True + with patch("weko_search_ui.tasks.get_redis_cache", return_value=None), \ + patch("weko_search_ui.tasks.RedisConnection") as redis_connection, \ + patch("weko_search_ui.tasks.inspect") as celery_inspect: + redis_connection.return_value.connection.return_value = datastore + + celery_inspect.return_value.ping.return_value = {"hostname": True} + assert check_celery_is_run(is_task=is_task) is True + datastore.put.assert_called_once_with(cache_key, b"1", 60) + + datastore.put.reset_mock() + celery_inspect.return_value.ping.return_value = {} + assert check_celery_is_run(is_task=is_task) is False + datastore.put.assert_called_once_with(cache_key, b"0", 60) + + datastore.put.reset_mock() + datastore.put.side_effect = RedisError("Redis error") + celery_inspect.return_value.ping.return_value = {"hostname": True} + with patch("weko_search_ui.tasks.current_app.logger.error") as mock_error: + assert check_celery_is_run(is_task=is_task) is True + datastore.put.assert_called_once_with(cache_key, b"1", 60) + assert "Redis error" in str(mock_error.call_args) + + with patch("weko_search_ui.tasks.get_redis_cache", return_value="1"): + assert check_celery_is_run(is_task=is_task) == True + with patch("weko_search_ui.tasks.get_redis_cache", return_value="0"): + assert check_celery_is_run(is_task=is_task) == False + + is_task = False + with patch("weko_search_ui.tasks.inspect") as celery_inspect: + celery_inspect.return_value.ping.return_value = {} + assert check_celery_is_run(is_task=is_task) is False class TestCheckSessionLifetime(unittest.TestCase): diff --git a/modules/weko-search-ui/weko_search_ui/admin.py b/modules/weko-search-ui/weko_search_ui/admin.py index f81b893e42..4c11956595 100644 --- a/modules/weko-search-ui/weko_search_ui/admin.py +++ b/modules/weko-search-ui/weko_search_ui/admin.py @@ -1194,7 +1194,7 @@ def check_export_status(self): """Check export status.""" if not current_user.is_authenticated: abort(302) - check_celery = check_celery_is_run() + check_celery = check_celery_is_run(is_task=True) check_life_time = check_session_lifetime() export_status, download_uri, message, run_message, \ status, start_time, finish_time = get_export_status() diff --git a/modules/weko-search-ui/weko_search_ui/config.py b/modules/weko-search-ui/weko_search_ui/config.py index fb159984a9..edf1bd1a43 100644 --- a/modules/weko-search-ui/weko_search_ui/config.py +++ b/modules/weko-search-ui/weko_search_ui/config.py @@ -759,6 +759,12 @@ WEKO_SEARCH_UI_BULK_EXPORT_RETRY_INTERVAL = 1 """ retry interval(sec) """ +WEKO_SEARCH_UI_CELERY_STATUS = "weko_search_ui_celery_status" +"""Cache key for storing the status of the Celery worker.""" + +WEKO_SEARCH_UI_CELERY_STATUS_CACHE_TTL = 60 +"""Celery status cache TTL in seconds.""" + WEKO_SEARCH_UI_IMPORT_REPLACE_RULES = {} """Strings to be replaced during item import.""" diff --git a/modules/weko-search-ui/weko_search_ui/tasks.py b/modules/weko-search-ui/weko_search_ui/tasks.py index 1007a5d148..c30d93d04b 100644 --- a/modules/weko-search-ui/weko_search_ui/tasks.py +++ b/modules/weko-search-ui/weko_search_ui/tasks.py @@ -38,6 +38,7 @@ from weko_admin.api import TempDirInfo from weko_admin.utils import get_redis_cache, reset_redis_cache from weko_redis.redis import RedisConnection +from redis import RedisError from invenio_db import db from .utils import ( @@ -315,12 +316,42 @@ def is_import_running(): return "is_import_running" -def check_celery_is_run(): +def check_celery_is_run(is_task=False): """Check celery is running, or not.""" - if not inspect(timeout=current_app.config.get("CELERY_GET_STATUS_TIMEOUT", 3.0)).ping(): - return False - else: - return True + cache_key = current_app.config.get("WEKO_SEARCH_UI_CELERY_STATUS", "weko_search_ui_celery_status") + cache_ttl = int(current_app.config.get( + "WEKO_SEARCH_UI_CELERY_STATUS_CACHE_TTL", 60 + )) + + cached_status = get_redis_cache(cache_key) + if cached_status and is_task: + if cached_status == "1": + return True + elif cached_status == "0": + return False + + is_running = bool( + inspect( + timeout=current_app.config.get("CELERY_GET_STATUS_TIMEOUT", 3.0) + ).ping() + ) + if is_task: + try: + redis_connection = RedisConnection() + datastore = redis_connection.connection( + db=current_app.config["CACHE_REDIS_DB"], kv=True + ) + datastore.put( + cache_key, + ("1" if is_running else "0").encode("utf-8"), + cache_ttl, + ) + except RedisError as ex: + current_app.logger.error( + "Could not cache Celery status; returning the ping result: %s", + ex, + ) + return is_running def check_session_lifetime(): """Check session lifetime."""