From 877609dfb6b9b64c767cc3879aa6ea61cf3c8a68 Mon Sep 17 00:00:00 2001 From: Alex Date: Sat, 19 Sep 2026 14:33:14 +0100 Subject: [PATCH] fix(worker): record worker state at startup, for every pool in_worker() read task_join_will_block, which eventlet and gevent leave unset, and the task's current_worker_task, which they scope to one greenlet -- so a greenlet a task spawned in those pools still took the web-process branch and dispatched to its own worker, where it could queue behind its parent and time out. The worker's own startup now records it: worker_init fires in every worker's main process before the pool starts (where solo, threads, eventlet and gevent run tasks, and what prefork children fork from), and worker_process_init in each prefork child. worker_ready would be too late -- prefork children are forked before it fires. The two existing checks stay for anything that runs tasks without that startup. --- docsgpt/celery_init.py | 45 +++++++++++++++++++++++++++++++----------- tests/test_celery.py | 22 +++++++++++++++++++++ 2 files changed, 56 insertions(+), 11 deletions(-) diff --git a/docsgpt/celery_init.py b/docsgpt/celery_init.py index 5e11d3e6..2eda5a3a 100644 --- a/docsgpt/celery_init.py +++ b/docsgpt/celery_init.py @@ -13,6 +13,7 @@ from celery.signals import ( setup_logging, task_postrun, task_prerun, + worker_init, worker_process_init, worker_ready, ) @@ -173,27 +174,49 @@ celery = make_celery() celery.config_from_object("docsgpt.celeryconfig") +#: Set once this process starts as a worker; see :func:`_mark_worker_process`. +_IS_WORKER_PROCESS = False + + +@worker_init.connect +@worker_process_init.connect +def _mark_worker_process(*args, **kwargs): + """Record that this process runs tasks, for :func:`in_worker`. + + ``worker_init`` fires in every worker's main process before its pool + starts: that is where solo, threads, eventlet and gevent run tasks, and + what prefork children fork from. ``worker_process_init`` covers prefork + children however they were started. + """ + global _IS_WORKER_PROCESS + _IS_WORKER_PROCESS = True + + def in_worker() -> bool: - """True anywhere in a Celery worker process, on any thread. + """True anywhere in a Celery worker process, on any thread or greenlet. ``current_worker_task`` alone is not enough: Celery records the executing - task on the thread that runs it, so a thread the task starts sees none and - would take the web-process branch — dispatching to the worker it is running - in and blocking on the result. Celery refuses that ``get()`` ("Never call - result.get() within a task!"), or, where joins are allowed, it waits on a - queue only this busy process serves. + task on the thread (or greenlet) that runs it, so one the task starts sees + none and would take the web-process branch — dispatching to the worker it + is running in and blocking on the result. Celery refuses that ``get()`` + ("Never call result.get() within a task!"), or, where joins are allowed, + it waits on a queue that only this busy process may be able to serve. - ``task_join_will_block`` is process-wide and set for every blocking pool - (prefork, solo, threads) — exactly the condition under which dispatching - and waiting goes wrong. eventlet/gevent leave it unset, so the task's own - thread still counts through ``current_worker_task``. + The worker's own startup (:func:`_mark_worker_process`) answers for every + pool. ``task_join_will_block`` — process-wide, set for every blocking pool + — and the task's own ``current_worker_task`` still count for a process + that runs tasks without having gone through that startup. Returns: bool: Whether this call is running inside a worker process. """ from celery.result import task_join_will_block - return task_join_will_block() or celery.current_worker_task is not None + return ( + _IS_WORKER_PROCESS + or task_join_will_block() + or celery.current_worker_task is not None + ) #: Task-name prefix the package carried before the rename to ``docsgpt``. diff --git a/tests/test_celery.py b/tests/test_celery.py index 5a3e66be..897b2623 100644 --- a/tests/test_celery.py +++ b/tests/test_celery.py @@ -317,6 +317,28 @@ class TestInWorker: with denied_join_result(): assert self._ask_from_a_new_thread() is True + def test_true_on_any_thread_of_a_non_blocking_pool_worker(self, monkeypatch): + # eventlet/gevent leave the join flag unset and scope the current task + # to one greenlet, so only the worker's own startup can say this + # process is a worker. The lifecycle signal records that for every + # thread and greenlet in it. + import docsgpt.celery_init as celery_init + + monkeypatch.setattr(celery_init, "_IS_WORKER_PROCESS", False) + assert self._ask_from_a_new_thread() is False + + celery_init.worker_init.send(sender=None) + + assert self._ask_from_a_new_thread() is True + + def test_prefork_children_record_it_on_their_own_start(self, monkeypatch): + import docsgpt.celery_init as celery_init + + monkeypatch.setattr(celery_init, "_IS_WORKER_PROCESS", False) + celery_init.worker_process_init.send(sender=None) + + assert self._ask_from_a_new_thread() is True + def test_true_on_the_task_thread_of_a_non_blocking_pool(self): # eventlet/gevent pools leave the process flag unset; the thread # running the task still knows it is in one.