fix: retain handles + done_callback + shutdown-cancel for pool cleanup tasks

Issue #25216 named two root causes for unbounded SESSION_POOL/USAGE_POOL
growth. The first — `periodic_usage_pool_cleanup` exiting permanently on
Redis lock acquire/renew failure — is being addressed by PR #21798's
`run_with_lock` helper. This commit closes the second: the
`asyncio.create_task(...)` calls at the cleanup-daemon spawn sites
discard their return value, so if either coroutine ever does exit for a
reason that escapes the inner loop (notably `asyncio.CancelledError`,
which is `BaseException` and is not caught by the `except Exception`
guards in `run_with_lock`), there is no real-time operator signal — the
asyncio default exception handler emits only a delayed "Task exception
was never retrieved" WARNING at GC time. By that point both pools have
already been growing unbounded for an unknown interval.

Three changes against `backend/open_webui/main.py`, mirroring the
existing `app.state.redis_task_command_listener` pattern at lines 676
and 754:

1. Store each cleanup-task handle on `app.state` and give it a `name=`
   so log lines identify which daemon died.
2. Attach a `_daemon_done_cb` that logs the exception at `ERROR` with
   `exc_info` at the moment of task death, not at GC time. The callback
   filters out `task.cancelled()` so it produces no noise during normal
   shutdown.
3. Cancel both tasks cooperatively in the lifespan shutdown block.
   Without this, `CancelledError` is never delivered, the tasks die
   while pending, and asyncio emits "Task was destroyed but it is
   pending!" warnings on every clean restart — which trains operators
   to ignore that warning and masks real failures.

The change is orthogonal to PR #21798 (zero file overlap; that PR
modifies `socket/main.py` and `socket/utils.py`, this one touches only
`main.py`). The two PRs can be reviewed, merged, and reverted
independently. With both in place, the cleanup loop is internally
resilient (it will not die from Redis events) AND the spawn site is
observable (if it dies anyway, the operator knows immediately).

Note: line 687 spawns `scheduler_worker_loop` with the same
fire-and-forget pattern. It is out of scope for this PR — issue #25216
specifically tracks the pool cleanup daemons — but the same treatment
would apply.

Refs: #25216, #21798
This commit is contained in:
Nexory 2026-05-30 17:34:04 +02:00
parent 4923920bf1
commit 19222eff68

View file

@ -660,8 +660,31 @@ async def lifespan(app: FastAPI):
limiter = anyio.to_thread.current_default_thread_limiter()
limiter.total_tokens = THREAD_POOL_SIZE
asyncio.create_task(periodic_usage_pool_cleanup())
asyncio.create_task(periodic_session_pool_cleanup())
def _daemon_done_cb(task: asyncio.Task) -> None:
# `asyncio.create_task` discards exceptions on unreferenced tasks; without
# this callback a dead cleanup loop only surfaces as a delayed
# "Task exception was never retrieved" WARNING at GC time, by which point
# SESSION_POOL/USAGE_POOL have already been growing without bound for an
# unknown interval. Logging at ERROR with exc_info gives operators an
# immediate, alertable signal at the moment of death.
if not task.cancelled() and task.exception() is not None:
log.error(
'Background daemon task %s exited unexpectedly: %r '
'— pool cleanup has stopped for the remaining lifetime of the process.',
task.get_name(),
task.exception(),
exc_info=task.exception(),
)
app.state.periodic_usage_cleanup_task = asyncio.create_task(
periodic_usage_pool_cleanup(), name='periodic_usage_pool_cleanup'
)
app.state.periodic_usage_cleanup_task.add_done_callback(_daemon_done_cb)
app.state.periodic_session_cleanup_task = asyncio.create_task(
periodic_session_pool_cleanup(), name='periodic_session_pool_cleanup'
)
app.state.periodic_session_cleanup_task.add_done_callback(_daemon_done_cb)
from open_webui.utils.automations import scheduler_worker_loop
@ -735,6 +758,16 @@ async def lifespan(app: FastAPI):
if hasattr(app.state, 'redis_task_command_listener'):
app.state.redis_task_command_listener.cancel()
# Cooperatively cancel the cleanup daemons so their `finally: release_fn()`
# blocks run and the Redis lock is freed for the next replica. Without this,
# the tasks are destroyed mid-await and asyncio emits a
# "Task was destroyed but it is pending!" warning on every clean restart,
# which trains operators to ignore that warning and masks real failures.
for _attr in ('periodic_usage_cleanup_task', 'periodic_session_cleanup_task'):
_t = getattr(app.state, _attr, None)
if _t is not None:
_t.cancel()
app = FastAPI(
title='Open WebUI',