From 7d8ec6fdbc614f6a49d2c79a161a1b640df406db Mon Sep 17 00:00:00 2001 From: datenzauberai Date: Fri, 10 Apr 2026 10:00:47 +0200 Subject: [PATCH 1/8] Preparing tests for additional celery backend --- tests/async_tests/conftest.py | 2 +- tests/async_tests/utils.py | 6 +++--- tests/background_callback/conftest.py | 2 +- tests/background_callback/utils.py | 6 +++--- 4 files changed, 8 insertions(+), 8 deletions(-) diff --git a/tests/async_tests/conftest.py b/tests/async_tests/conftest.py index b701eea91a..1f225d322d 100644 --- a/tests/async_tests/conftest.py +++ b/tests/async_tests/conftest.py @@ -4,7 +4,7 @@ if "REDIS_URL" in os.environ: - managers = ["celery", "diskcache"] + managers = ["celery-redis", "diskcache"] else: print("Skipping celery tests because REDIS_URL is not defined") managers = ["diskcache"] diff --git a/tests/async_tests/utils.py b/tests/async_tests/utils.py index b7074b0735..2b88f8a3ca 100644 --- a/tests/async_tests/utils.py +++ b/tests/async_tests/utils.py @@ -36,7 +36,7 @@ def get_background_callback_manager(): """ Get the long callback mangaer configured by environment variables """ - if os.environ.get("LONG_CALLBACK_MANAGER", None) == "celery": + if os.environ.get("LONG_CALLBACK_MANAGER", None) == "celery-redis": from dash.background_callback import CeleryManager from celery import Celery @@ -77,8 +77,8 @@ def kill(proc_pid): def setup_background_callback_app(manager_name, app_name): from dash.testing.application_runners import import_app - if manager_name == "celery": - os.environ["LONG_CALLBACK_MANAGER"] = "celery" + if manager_name == "celery-redis": + os.environ["LONG_CALLBACK_MANAGER"] = "celery-redis" redis_url = os.environ["REDIS_URL"].rstrip("/") os.environ["CELERY_BROKER"] = f"{redis_url}/0" os.environ["CELERY_BACKEND"] = f"{redis_url}/1" diff --git a/tests/background_callback/conftest.py b/tests/background_callback/conftest.py index b701eea91a..1f225d322d 100644 --- a/tests/background_callback/conftest.py +++ b/tests/background_callback/conftest.py @@ -4,7 +4,7 @@ if "REDIS_URL" in os.environ: - managers = ["celery", "diskcache"] + managers = ["celery-redis", "diskcache"] else: print("Skipping celery tests because REDIS_URL is not defined") managers = ["diskcache"] diff --git a/tests/background_callback/utils.py b/tests/background_callback/utils.py index 1cefd4ecc3..60c35487d2 100644 --- a/tests/background_callback/utils.py +++ b/tests/background_callback/utils.py @@ -35,7 +35,7 @@ def get_background_callback_manager(): """ Get the long callback mangaer configured by environment variables """ - if os.environ.get("LONG_CALLBACK_MANAGER", None) == "celery": + if os.environ.get("LONG_CALLBACK_MANAGER", None) == "celery-redis": from dash.background_callback import CeleryManager from celery import Celery import redis @@ -77,8 +77,8 @@ def kill(proc_pid): def setup_background_callback_app(manager_name, app_name): from dash.testing.application_runners import import_app - if manager_name == "celery": - os.environ["LONG_CALLBACK_MANAGER"] = "celery" + if manager_name == "celery-redis": + os.environ["LONG_CALLBACK_MANAGER"] = "celery-redis" redis_url = os.environ["REDIS_URL"].rstrip("/") os.environ["CELERY_BROKER"] = f"{redis_url}/0" os.environ["CELERY_BACKEND"] = f"{redis_url}/1" From 808d0c9745450b736ed33550b4c331a09883dabc Mon Sep 17 00:00:00 2001 From: datenzauberai Date: Mon, 3 Aug 2026 13:04:38 +0200 Subject: [PATCH 2/8] Draft implementation to support celery-filesystem provider --- .../managers/celery_manager.py | 74 +++++++++++++----- tests/background_callback/conftest.py | 5 +- .../test_basic_long_callback001.py | 2 +- tests/background_callback/utils.py | 75 +++++++++++++------ 4 files changed, 115 insertions(+), 41 deletions(-) diff --git a/dash/background_callback/managers/celery_manager.py b/dash/background_callback/managers/celery_manager.py index eb2864caa0..eb1f7916fa 100644 --- a/dash/background_callback/managers/celery_manager.py +++ b/dash/background_callback/managers/celery_manager.py @@ -1,5 +1,6 @@ import inspect import json +import logging import traceback from contextvars import copy_context import asyncio @@ -13,6 +14,12 @@ from dash.background_callback._proxy_set_props import ProxySetProps from dash.background_callback.managers import BaseBackgroundCallbackManager +logging.basicConfig( + encoding="utf-8", level=logging.DEBUG, format="%(levelname)s:%(name)s:%(message)s" +) +logging.getLogger("selenium.webdriver.remote.remote_connection").setLevel(logging.INFO) +logger = logging.getLogger(__name__) + class CeleryManager(BaseBackgroundCallbackManager): """Manage background execution of callbacks with a celery queue.""" @@ -41,12 +48,10 @@ def __init__(self, celery_app, cache_by=None, expire=None): DisabledBackend, ) except ImportError as missing_imports: - raise ImportError( - """\ + raise ImportError("""\ CeleryManager requires extra dependencies which can be installed doing - $ pip install "dash[celery]"\n""" - ) from missing_imports + $ pip install "dash[celery]"\n""") from missing_imports if not isinstance(celery_app, celery.Celery): raise ValueError("First argument must be a celery.Celery object") @@ -67,6 +72,7 @@ def terminate_job(self, job): def terminate_unhealthy_job(self, job): task = self.get_task(job) if task and task.status in ("FAILURE", "REVOKED"): + logger.info("Terminating unhealthy job %s with status %s", job, task.status) return self.terminate_job(job) return False @@ -89,8 +95,18 @@ def get_task(self, job): return None + @staticmethod + def _ensure_bytes(o) -> bytes: + if isinstance(o, bytes): + return o + return str(o).encode() + def clear_cache_entry(self, key): - self.handle.backend.delete(key) + logger.info("Clearing cache entry for %s", key) + # delete should not be called when the entry is not present + value = self.handle.backend.get(self._ensure_bytes(key)) + if value is not None: + self.handle.backend.delete(self._ensure_bytes(key)) def get_or_create_signing_secret(self, generate): backend = self.handle.backend @@ -103,24 +119,33 @@ def get_or_create_signing_secret(self, generate): return backend.get(self.SIGNING_SECRET_KEY) or secret def call_job_fn(self, key, job_fn, args, context): - task = job_fn.delay(key, self._make_progress_key(key), args, context) + result_key = self._ensure_bytes(key) + progress_key = self._ensure_bytes(self._make_progress_key(key)) + set_props_key = self._ensure_bytes(self._make_set_props_key(key)) + task = job_fn.delay(result_key, progress_key, set_props_key, args, context) return task.task_id def get_progress(self, key): - progress_key = self._make_progress_key(key) + progress_key = self._ensure_bytes(self._make_progress_key(key)) + logger.info("Getting progress for %s", progress_key) progress_data = self.handle.backend.get(progress_key) if progress_data: - self.handle.backend.delete(progress_key) + self.clear_cache_entry(progress_key) return json.loads(progress_data) return None def result_ready(self, key): - return self.handle.backend.get(key) is not None + result_key = self._ensure_bytes(key) + logger.info("Getting result for %s", result_key) + return self.handle.backend.get(result_key) is not None def get_result(self, key, job): + result_key = self._ensure_bytes(key) + progress_key = self._ensure_bytes(self._make_progress_key(key)) # Get result value - result = self.handle.backend.get(key) + logger.info("Getting result for %s", key) + result = self.handle.backend.get(result_key) if result is None: return self.UNDEFINED @@ -128,22 +153,24 @@ def get_result(self, key, job): # Clear result if not caching if self.cache_by is None: - self.clear_cache_entry(key) + self.clear_cache_entry(result_key) else: if self.expire: # Set/update expiration time - self.handle.backend.expire(key, self.expire) - self.clear_cache_entry(self._make_progress_key(key)) + self.handle.backend.expire(result_key, self.expire) + self.clear_cache_entry(progress_key) self.terminate_job(job) return result def get_updated_props(self, key): - updated_props = self.handle.backend.get(self._make_set_props_key(key)) + set_props_key = self._ensure_bytes(self._make_set_props_key(key)) + logger.info("Getting updated props for %s", set_props_key) + updated_props = self.handle.backend.get(set_props_key) if updated_props is None: return {} - self.clear_cache_entry(key) + self.clear_cache_entry(set_props_key) return json.loads(updated_props) @@ -153,25 +180,28 @@ def _make_job_fn(fn, celery_app, progress, key): # pylint: disable=too-many-sta @celery_app.task(name=f"background_callback_{key}") def job_fn( - result_key, progress_key, user_callback_args, context=None + result_key, progress_key, set_props_key, user_callback_args, context=None ): # pylint: disable=too-many-statements def _set_progress(progress_value): if not isinstance(progress_value, (list, tuple)): progress_value = [progress_value] + logger.info("Setting progress for %s to %s", progress_key, progress_value) cache.set(progress_key, json.dumps(progress_value, cls=PlotlyJSONEncoder)) maybe_progress = [_set_progress] if progress else [] def _set_props(_id, props): + logger.info("Setting updated props for %s to %s", set_props_key, props) cache.set( - f"{result_key}-set_props", + set_props_key, json.dumps({_id: props}, cls=PlotlyJSONEncoder), ) ctx = copy_context() def run(): + logger.info("Running callback for %s", result_key) c = AttributeDict(**context) # type: ignore[reportCallIssue] c.ignore_register_page = False c.updated_props = ProxySetProps(_set_props) @@ -188,6 +218,7 @@ def run(): except PreventUpdate: # Put NoUpdate dict directly to avoid circular imports. errored = True + logger.info("Callback prevented update for %s", result_key) cache.set( result_key, json.dumps( @@ -195,6 +226,7 @@ def run(): ), ) except Exception as err: # pylint: disable=broad-except + logger.info("Callback failed with error: %s", err) errored = True cache.set( result_key, @@ -209,6 +241,9 @@ def run(): ) if not errored: + logger.info( + "Setting result for %s to %s", result_key, user_callback_output + ) cache.set( result_key, json.dumps(user_callback_output, cls=PlotlyJSONEncoder) ) @@ -234,6 +269,7 @@ async def async_run(): except PreventUpdate: # Put NoUpdate dict directly to avoid circular imports. errored = True + logger.info("Callback prevented update for %s", result_key) cache.set( result_key, json.dumps( @@ -242,6 +278,7 @@ async def async_run(): ) except Exception as err: # pylint: disable=broad-except errored = True + logger.info("Callback failed with error: %s", err) cache.set( result_key, json.dumps( @@ -258,6 +295,9 @@ async def async_run(): if asyncio.iscoroutine(user_callback_output): user_callback_output = await user_callback_output + logger.info( + "Setting result for %s to %s", result_key, user_callback_output + ) cache.set( result_key, json.dumps(user_callback_output, cls=PlotlyJSONEncoder) ) diff --git a/tests/background_callback/conftest.py b/tests/background_callback/conftest.py index 1f225d322d..125d0fe15e 100644 --- a/tests/background_callback/conftest.py +++ b/tests/background_callback/conftest.py @@ -2,12 +2,13 @@ import pytest +os.environ["REDIS_URL"] = "redis://localhost:6379" if "REDIS_URL" in os.environ: - managers = ["celery-redis", "diskcache"] + managers = ["celery-filesystem", "celery-redis", "diskcache"] else: print("Skipping celery tests because REDIS_URL is not defined") - managers = ["diskcache"] + managers = ["celery-filesystem", "diskcache"] @pytest.fixture(params=managers) diff --git a/tests/background_callback/test_basic_long_callback001.py b/tests/background_callback/test_basic_long_callback001.py index af8c80d3ee..0f3fa31026 100644 --- a/tests/background_callback/test_basic_long_callback001.py +++ b/tests/background_callback/test_basic_long_callback001.py @@ -10,7 +10,7 @@ @pytest.mark.skipif( sys.version_info < (3, 7), reason="Python 3.6 long callbacks tests hangs up" ) -@flaky(max_runs=3) +@flaky(max_runs=1) def test_lcbc001_fast_input(dash_duo, manager): """ Make sure that we settle to the correct final value when handling rapid inputs diff --git a/tests/background_callback/utils.py b/tests/background_callback/utils.py index 60c35487d2..fe9bc6887e 100644 --- a/tests/background_callback/utils.py +++ b/tests/background_callback/utils.py @@ -1,3 +1,4 @@ +import logging import os import sys import shutil @@ -48,6 +49,24 @@ def get_background_callback_manager(): background_callback_manager = CeleryManager(celery_app) redis_conn = redis.Redis(host="localhost", port=6379, db=1) background_callback_manager.test_lock = redis_conn.lock("test-lock") + elif os.environ.get("LONG_CALLBACK_MANAGER", None) == "celery-filesystem": + from dash.background_callback import CeleryManager + from celery import Celery + + celery_app = Celery( + __name__, + broker=os.environ.get("CELERY_BROKER"), + backend=os.environ.get("CELERY_BACKEND"), + broker_transport_options={ + "data_folder_in": os.environ.get("CELERY_BROKER_FILESYSTEM_DIRECTORY"), + "data_folder_out": os.environ.get("CELERY_BROKER_FILESYSTEM_DIRECTORY"), + "control_folder": os.environ.get("CELERY_BROKER_FILESYSTEM_DIRECTORY"), + }, + ) + background_callback_manager = CeleryManager(celery_app) + + # TODO implement lock based on filesystem for testing? + # background_callback_manager.test_lock = ??? elif os.environ.get("LONG_CALLBACK_MANAGER", None) == "diskcache": import diskcache @@ -77,17 +96,28 @@ def kill(proc_pid): def setup_background_callback_app(manager_name, app_name): from dash.testing.application_runners import import_app - if manager_name == "celery-redis": - os.environ["LONG_CALLBACK_MANAGER"] = "celery-redis" - redis_url = os.environ["REDIS_URL"].rstrip("/") - os.environ["CELERY_BROKER"] = f"{redis_url}/0" - os.environ["CELERY_BACKEND"] = f"{redis_url}/1" - - # Clear redis of cached values - redis_conn = redis.Redis(host="localhost", port=6379, db=1) - cache_keys = redis_conn.keys() - if cache_keys: - redis_conn.delete(*cache_keys) + if manager_name in ["celery-redis", "celery-filesystem"]: + os.environ["LONG_CALLBACK_MANAGER"] = manager_name + + if manager_name == "celery-redis": + redis_url = os.environ["REDIS_URL"].rstrip("/") + os.environ["CELERY_BROKER"] = f"{redis_url}/0" + os.environ["CELERY_BACKEND"] = f"{redis_url}/1" + + # Clear redis of cached values + redis_conn = redis.Redis(host="localhost", port=6379, db=1) + cache_keys = redis_conn.keys() + if cache_keys: + redis_conn.delete(*cache_keys) + elif manager_name == "celery-filesystem": + # celery_filesystem_directory = tempfile.mkdtemp(prefix="lc-celery-") + celery_filesystem_directory = "/tmp/lc-celery-broker-filesystem" + os.environ["CELERY_BROKER"] = "filesystem://" + os.environ["CELERY_BROKER_FILESYSTEM_DIRECTORY"] = ( + celery_filesystem_directory + ) + print(f"{celery_filesystem_directory=}") + os.environ["CELERY_BACKEND"] = f"file://{celery_filesystem_directory}" worker = subprocess.Popen( [ @@ -102,22 +132,25 @@ def setup_background_callback_app(manager_name, app_name): "--concurrency", "2", "--loglevel=info", + "--logfile=/tmp/lc-celery-broker-filesystem/celery_worker_%i.log", ], encoding="utf8", preexec_fn=os.setpgrp, - stderr=subprocess.PIPE, + # stderr=subprocess.PIPE, ) + logging.debug(f"Started celery worker with PID {worker.pid}") # Wait for the worker to be ready, if you cancel before it is ready, the job # will still be queued. - lines = [] - for line in iter(worker.stderr.readline, ""): - if "ready" in line: - break - lines.append(line) - else: - error = "\n".join(lines) - error += f"\nPath: {sys.path}" - raise RuntimeError(f"celery failed to start: {error}") + time.sleep(5) + # lines = [] + # for line in iter(worker.stderr.readline, ""): + # if "ready" in line: + # break + # lines.append(line) + # else: + # error = "\n".join(lines) + # error += f"\nPath: {sys.path}" + # raise RuntimeError(f"celery failed to start: {error}") try: yield import_app(f"tests.background_callback.{app_name}") From aec4a3f74cac42a83631fe2a787283c833dafea9 Mon Sep 17 00:00:00 2001 From: datenzauberai Date: Mon, 3 Aug 2026 16:36:39 +0200 Subject: [PATCH 3/8] Use filelock as a test-lock for the celery-filesystem backend --- requirements/ci.txt | 1 + tests/background_callback/utils.py | 17 ++++++++++++----- 2 files changed, 13 insertions(+), 5 deletions(-) diff --git a/requirements/ci.txt b/requirements/ci.txt index 8e18280d04..2581571520 100644 --- a/requirements/ci.txt +++ b/requirements/ci.txt @@ -2,6 +2,7 @@ black==22.3.0 flake8==7.0.0 flaky==3.8.1 +filelock>=3.0 flask-talisman==1.0.0 ipython<9.0.0 mimesis<=11.1.0 diff --git a/tests/background_callback/utils.py b/tests/background_callback/utils.py index fe9bc6887e..1ca6948c74 100644 --- a/tests/background_callback/utils.py +++ b/tests/background_callback/utils.py @@ -52,21 +52,28 @@ def get_background_callback_manager(): elif os.environ.get("LONG_CALLBACK_MANAGER", None) == "celery-filesystem": from dash.background_callback import CeleryManager from celery import Celery + from filelock import FileLock + + celery_broker_path = os.environ.get("CELERY_BROKER_FILESYSTEM_DIRECTORY") + assert ( + celery_broker_path is not None + ), "CELERY_BROKER_FILESYSTEM_DIRECTORY must be set" celery_app = Celery( __name__, broker=os.environ.get("CELERY_BROKER"), backend=os.environ.get("CELERY_BACKEND"), broker_transport_options={ - "data_folder_in": os.environ.get("CELERY_BROKER_FILESYSTEM_DIRECTORY"), - "data_folder_out": os.environ.get("CELERY_BROKER_FILESYSTEM_DIRECTORY"), - "control_folder": os.environ.get("CELERY_BROKER_FILESYSTEM_DIRECTORY"), + "data_folder_in": celery_broker_path, + "data_folder_out": celery_broker_path, + "control_folder": celery_broker_path, }, ) background_callback_manager = CeleryManager(celery_app) - # TODO implement lock based on filesystem for testing? - # background_callback_manager.test_lock = ??? + background_callback_manager.test_lock = FileLock( + os.path.join(celery_broker_path, "test-lock") + ) elif os.environ.get("LONG_CALLBACK_MANAGER", None) == "diskcache": import diskcache From 8a46d7e8005bc037185253f725f618e0cf65fe17 Mon Sep 17 00:00:00 2001 From: datenzauberai Date: Mon, 3 Aug 2026 16:49:15 +0200 Subject: [PATCH 4/8] Remvoe debug logging --- .../managers/celery_manager.py | 26 ------------------ .../test_basic_long_callback001.py | 2 +- tests/background_callback/utils.py | 27 ++++++++----------- 3 files changed, 12 insertions(+), 43 deletions(-) diff --git a/dash/background_callback/managers/celery_manager.py b/dash/background_callback/managers/celery_manager.py index eb1f7916fa..fe5b84159b 100644 --- a/dash/background_callback/managers/celery_manager.py +++ b/dash/background_callback/managers/celery_manager.py @@ -1,6 +1,5 @@ import inspect import json -import logging import traceback from contextvars import copy_context import asyncio @@ -14,12 +13,6 @@ from dash.background_callback._proxy_set_props import ProxySetProps from dash.background_callback.managers import BaseBackgroundCallbackManager -logging.basicConfig( - encoding="utf-8", level=logging.DEBUG, format="%(levelname)s:%(name)s:%(message)s" -) -logging.getLogger("selenium.webdriver.remote.remote_connection").setLevel(logging.INFO) -logger = logging.getLogger(__name__) - class CeleryManager(BaseBackgroundCallbackManager): """Manage background execution of callbacks with a celery queue.""" @@ -72,7 +65,6 @@ def terminate_job(self, job): def terminate_unhealthy_job(self, job): task = self.get_task(job) if task and task.status in ("FAILURE", "REVOKED"): - logger.info("Terminating unhealthy job %s with status %s", job, task.status) return self.terminate_job(job) return False @@ -102,7 +94,6 @@ def _ensure_bytes(o) -> bytes: return str(o).encode() def clear_cache_entry(self, key): - logger.info("Clearing cache entry for %s", key) # delete should not be called when the entry is not present value = self.handle.backend.get(self._ensure_bytes(key)) if value is not None: @@ -127,7 +118,6 @@ def call_job_fn(self, key, job_fn, args, context): def get_progress(self, key): progress_key = self._ensure_bytes(self._make_progress_key(key)) - logger.info("Getting progress for %s", progress_key) progress_data = self.handle.backend.get(progress_key) if progress_data: self.clear_cache_entry(progress_key) @@ -137,14 +127,12 @@ def get_progress(self, key): def result_ready(self, key): result_key = self._ensure_bytes(key) - logger.info("Getting result for %s", result_key) return self.handle.backend.get(result_key) is not None def get_result(self, key, job): result_key = self._ensure_bytes(key) progress_key = self._ensure_bytes(self._make_progress_key(key)) # Get result value - logger.info("Getting result for %s", key) result = self.handle.backend.get(result_key) if result is None: return self.UNDEFINED @@ -165,7 +153,6 @@ def get_result(self, key, job): def get_updated_props(self, key): set_props_key = self._ensure_bytes(self._make_set_props_key(key)) - logger.info("Getting updated props for %s", set_props_key) updated_props = self.handle.backend.get(set_props_key) if updated_props is None: return {} @@ -186,13 +173,11 @@ def _set_progress(progress_value): if not isinstance(progress_value, (list, tuple)): progress_value = [progress_value] - logger.info("Setting progress for %s to %s", progress_key, progress_value) cache.set(progress_key, json.dumps(progress_value, cls=PlotlyJSONEncoder)) maybe_progress = [_set_progress] if progress else [] def _set_props(_id, props): - logger.info("Setting updated props for %s to %s", set_props_key, props) cache.set( set_props_key, json.dumps({_id: props}, cls=PlotlyJSONEncoder), @@ -201,7 +186,6 @@ def _set_props(_id, props): ctx = copy_context() def run(): - logger.info("Running callback for %s", result_key) c = AttributeDict(**context) # type: ignore[reportCallIssue] c.ignore_register_page = False c.updated_props = ProxySetProps(_set_props) @@ -218,7 +202,6 @@ def run(): except PreventUpdate: # Put NoUpdate dict directly to avoid circular imports. errored = True - logger.info("Callback prevented update for %s", result_key) cache.set( result_key, json.dumps( @@ -226,7 +209,6 @@ def run(): ), ) except Exception as err: # pylint: disable=broad-except - logger.info("Callback failed with error: %s", err) errored = True cache.set( result_key, @@ -241,9 +223,6 @@ def run(): ) if not errored: - logger.info( - "Setting result for %s to %s", result_key, user_callback_output - ) cache.set( result_key, json.dumps(user_callback_output, cls=PlotlyJSONEncoder) ) @@ -269,7 +248,6 @@ async def async_run(): except PreventUpdate: # Put NoUpdate dict directly to avoid circular imports. errored = True - logger.info("Callback prevented update for %s", result_key) cache.set( result_key, json.dumps( @@ -278,7 +256,6 @@ async def async_run(): ) except Exception as err: # pylint: disable=broad-except errored = True - logger.info("Callback failed with error: %s", err) cache.set( result_key, json.dumps( @@ -295,9 +272,6 @@ async def async_run(): if asyncio.iscoroutine(user_callback_output): user_callback_output = await user_callback_output - logger.info( - "Setting result for %s to %s", result_key, user_callback_output - ) cache.set( result_key, json.dumps(user_callback_output, cls=PlotlyJSONEncoder) ) diff --git a/tests/background_callback/test_basic_long_callback001.py b/tests/background_callback/test_basic_long_callback001.py index 0f3fa31026..af8c80d3ee 100644 --- a/tests/background_callback/test_basic_long_callback001.py +++ b/tests/background_callback/test_basic_long_callback001.py @@ -10,7 +10,7 @@ @pytest.mark.skipif( sys.version_info < (3, 7), reason="Python 3.6 long callbacks tests hangs up" ) -@flaky(max_runs=1) +@flaky(max_runs=3) def test_lcbc001_fast_input(dash_duo, manager): """ Make sure that we settle to the correct final value when handling rapid inputs diff --git a/tests/background_callback/utils.py b/tests/background_callback/utils.py index 1ca6948c74..485ca11396 100644 --- a/tests/background_callback/utils.py +++ b/tests/background_callback/utils.py @@ -1,4 +1,3 @@ -import logging import os import sys import shutil @@ -117,8 +116,7 @@ def setup_background_callback_app(manager_name, app_name): if cache_keys: redis_conn.delete(*cache_keys) elif manager_name == "celery-filesystem": - # celery_filesystem_directory = tempfile.mkdtemp(prefix="lc-celery-") - celery_filesystem_directory = "/tmp/lc-celery-broker-filesystem" + celery_filesystem_directory = tempfile.mkdtemp(prefix="lc-celery-") os.environ["CELERY_BROKER"] = "filesystem://" os.environ["CELERY_BROKER_FILESYSTEM_DIRECTORY"] = ( celery_filesystem_directory @@ -139,25 +137,22 @@ def setup_background_callback_app(manager_name, app_name): "--concurrency", "2", "--loglevel=info", - "--logfile=/tmp/lc-celery-broker-filesystem/celery_worker_%i.log", ], encoding="utf8", preexec_fn=os.setpgrp, - # stderr=subprocess.PIPE, + stderr=subprocess.PIPE, ) - logging.debug(f"Started celery worker with PID {worker.pid}") # Wait for the worker to be ready, if you cancel before it is ready, the job # will still be queued. - time.sleep(5) - # lines = [] - # for line in iter(worker.stderr.readline, ""): - # if "ready" in line: - # break - # lines.append(line) - # else: - # error = "\n".join(lines) - # error += f"\nPath: {sys.path}" - # raise RuntimeError(f"celery failed to start: {error}") + lines = [] + for line in iter(worker.stderr.readline, ""): + if "ready" in line: + break + lines.append(line) + else: + error = "\n".join(lines) + error += f"\nPath: {sys.path}" + raise RuntimeError(f"celery failed to start: {error}") try: yield import_app(f"tests.background_callback.{app_name}") From 5d77ba944553335a44dc7fccfe816cb8cb65acba Mon Sep 17 00:00:00 2001 From: datenzauberai Date: Mon, 3 Aug 2026 17:24:57 +0200 Subject: [PATCH 5/8] Error early on when configure backend is not a BaseKeyValueStoreBackend --- dash/background_callback/managers/celery_manager.py | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/dash/background_callback/managers/celery_manager.py b/dash/background_callback/managers/celery_manager.py index fe5b84159b..a7bdb65683 100644 --- a/dash/background_callback/managers/celery_manager.py +++ b/dash/background_callback/managers/celery_manager.py @@ -39,6 +39,7 @@ def __init__(self, celery_app, cache_by=None, expire=None): import celery # type: ignore[import-not-found] # pylint: disable=import-outside-toplevel,import-error from celery.backends.base import ( # type: ignore[import-not-found] # pylint: disable=import-outside-toplevel,import-error DisabledBackend, + BaseKeyValueStoreBackend, ) except ImportError as missing_imports: raise ImportError("""\ @@ -52,6 +53,11 @@ def __init__(self, celery_app, cache_by=None, expire=None): if isinstance(celery_app.backend, DisabledBackend): raise ValueError("Celery instance must be configured with a result backend") + if not isinstance(celery_app.backend, BaseKeyValueStoreBackend): + raise ValueError( + "Celery must be configured with a key-value store backend (e.g. Redis or Filesystem)" + ) + self.handle = celery_app self.expire = expire super().__init__(cache_by) From f1fe60ebdf1d5b6e9373b9380abc306b0d89fe20 Mon Sep 17 00:00:00 2001 From: datenzauberai Date: Mon, 3 Aug 2026 17:40:12 +0200 Subject: [PATCH 6/8] Do not set redis url on test set-up. --- tests/background_callback/conftest.py | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/tests/background_callback/conftest.py b/tests/background_callback/conftest.py index 125d0fe15e..8ac581bd59 100644 --- a/tests/background_callback/conftest.py +++ b/tests/background_callback/conftest.py @@ -2,12 +2,10 @@ import pytest -os.environ["REDIS_URL"] = "redis://localhost:6379" - if "REDIS_URL" in os.environ: managers = ["celery-filesystem", "celery-redis", "diskcache"] else: - print("Skipping celery tests because REDIS_URL is not defined") + print("Skipping celery tests on Redis because REDIS_URL is not defined") managers = ["celery-filesystem", "diskcache"] From d5c26e4d6b73203d16b04f9d8326492eb235d17f Mon Sep 17 00:00:00 2001 From: datenzauberai Date: Tue, 4 Aug 2026 10:18:28 +0200 Subject: [PATCH 7/8] Add celery-filesystem manager to async_tests --- tests/async_tests/conftest.py | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/tests/async_tests/conftest.py b/tests/async_tests/conftest.py index 1f225d322d..8ac581bd59 100644 --- a/tests/async_tests/conftest.py +++ b/tests/async_tests/conftest.py @@ -2,12 +2,11 @@ import pytest - if "REDIS_URL" in os.environ: - managers = ["celery-redis", "diskcache"] + managers = ["celery-filesystem", "celery-redis", "diskcache"] else: - print("Skipping celery tests because REDIS_URL is not defined") - managers = ["diskcache"] + print("Skipping celery tests on Redis because REDIS_URL is not defined") + managers = ["celery-filesystem", "diskcache"] @pytest.fixture(params=managers) From 5dd4426db7733dfa0a478a489b7962e600448f02 Mon Sep 17 00:00:00 2001 From: datenzauberai Date: Tue, 4 Aug 2026 10:19:59 +0200 Subject: [PATCH 8/8] Avoid deprecation warning --- tests/background_callback/utils.py | 7 ++++--- 1 file changed, 4 insertions(+), 3 deletions(-) diff --git a/tests/background_callback/utils.py b/tests/background_callback/utils.py index 485ca11396..9beb37f000 100644 --- a/tests/background_callback/utils.py +++ b/tests/background_callback/utils.py @@ -44,6 +44,7 @@ def get_background_callback_manager(): __name__, broker=os.environ.get("CELERY_BROKER"), backend=os.environ.get("CELERY_BACKEND"), + broker_connection_retry_on_startup=True, ) background_callback_manager = CeleryManager(celery_app) redis_conn = redis.Redis(host="localhost", port=6379, db=1) @@ -118,9 +119,9 @@ def setup_background_callback_app(manager_name, app_name): elif manager_name == "celery-filesystem": celery_filesystem_directory = tempfile.mkdtemp(prefix="lc-celery-") os.environ["CELERY_BROKER"] = "filesystem://" - os.environ["CELERY_BROKER_FILESYSTEM_DIRECTORY"] = ( - celery_filesystem_directory - ) + os.environ[ + "CELERY_BROKER_FILESYSTEM_DIRECTORY" + ] = celery_filesystem_directory print(f"{celery_filesystem_directory=}") os.environ["CELERY_BACKEND"] = f"file://{celery_filesystem_directory}"