Merge pull request #22 from darth-veitcher/fix/run-async-real-comfyui-event-loop
fix: chat_structured() crashes in real ComfyUI's async execution engine
This commit is contained in:
@@ -88,17 +88,29 @@ _CHAT_RESPONSE_CACHE = _TTLLRUCache(maxsize=64, ttl_seconds=None)
|
||||
|
||||
|
||||
def _run_async(coro):
|
||||
"""Run an async coroutine synchronously, safe inside a running event loop."""
|
||||
try:
|
||||
asyncio.get_running_loop()
|
||||
# Called from within a running loop (e.g. ComfyUI's async executor).
|
||||
# Spin up a worker thread with its own loop to avoid "loop already running".
|
||||
import concurrent.futures
|
||||
"""Run an async coroutine synchronously in an isolated worker thread.
|
||||
|
||||
with concurrent.futures.ThreadPoolExecutor(max_workers=1) as pool:
|
||||
return pool.submit(asyncio.run, coro).result()
|
||||
except RuntimeError:
|
||||
return asyncio.run(coro)
|
||||
Always spins up a fresh thread rather than conditionally checking
|
||||
asyncio.get_running_loop() first: live-verified against a real running
|
||||
ComfyUI instance (its execution engine runs its own event loop, in
|
||||
Python 3.13, on the same process) that the conditional version — try
|
||||
get_running_loop(), spin up a thread only if it succeeds, otherwise
|
||||
call asyncio.run(coro) directly — is unreliable there. Under real
|
||||
ComfyUI, get_running_loop() sometimes raised inside that try block
|
||||
(unlike under pytest or a standalone script, where it never does),
|
||||
which routed straight into `asyncio.run(coro)` on the *current* thread
|
||||
— the one thread guaranteed to already have ComfyUI's own loop running
|
||||
— reproducing exactly the "asyncio.run() cannot be called from a
|
||||
running event loop" crash this function exists to prevent. Always
|
||||
using a dedicated thread sidesteps the detection entirely: a freshly
|
||||
spawned thread never has an ambient loop, so asyncio.run() is safe
|
||||
there unconditionally, regardless of what the calling thread's loop
|
||||
state actually is.
|
||||
"""
|
||||
import concurrent.futures
|
||||
|
||||
with concurrent.futures.ThreadPoolExecutor(max_workers=1) as pool:
|
||||
return pool.submit(asyncio.run, coro).result()
|
||||
|
||||
|
||||
async def _post_json(
|
||||
|
||||
@@ -12,6 +12,8 @@ BDD coverage:
|
||||
../specs/007-llm-provider-abstraction/features/us3_model_lifecycle.feature
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
|
||||
import pytest
|
||||
|
||||
import comfydv._llm.ollama_provider as provider_mod
|
||||
@@ -326,3 +328,71 @@ def test_chat_structured_caches_after_successful_validation(monkeypatch):
|
||||
|
||||
assert r1 == r2 == Widget(name="cached")
|
||||
assert calls["n"] == 1
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# _run_async — regression coverage for a bug found running the real
|
||||
# docker-compose ComfyUI harness (not caught by any prior test, and not
|
||||
# reproducible under pytest — see below). ComfyUI's actual async execution
|
||||
# engine runs node functions synchronously *inside* an already-running event
|
||||
# loop. The old implementation tried `asyncio.get_running_loop()` first and
|
||||
# only spun up an isolated worker thread if that succeeded — under real
|
||||
# ComfyUI (Python 3.13), that check sometimes raised anyway, falling through
|
||||
# to a direct `asyncio.run(coro)` call on the *current* thread — the one
|
||||
# thread guaranteed to already have a loop running — reproducing "asyncio.run()
|
||||
# cannot be called from a running event loop" exactly, every time
|
||||
# chat_structured() was called.
|
||||
#
|
||||
# Honest caveat (raised by beacon-reviewer, confirmed by reconstructing the
|
||||
# old code and running it against these tests): under CPython/pytest,
|
||||
# asyncio.get_running_loop() inside asyncio.run(outer()) never spuriously
|
||||
# raises, so the *old, buggy* implementation also takes the safe
|
||||
# worker-thread branch here and these tests pass against it too. They do not
|
||||
# reproduce the actual reported bug — only the live end-to-end run against
|
||||
# real ComfyUI did that. What these tests do lock in: the documented
|
||||
# contract ("calling _run_async from within a running loop must work and
|
||||
# must propagate exceptions"), so a regression to a naive
|
||||
# `asyncio.run(coro)`-with-no-thread implementation (which *does* fail here)
|
||||
# gets caught immediately.
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def test_run_async_works_with_no_ambient_loop():
|
||||
async def coro():
|
||||
return "ok"
|
||||
|
||||
assert _run_async(coro()) == "ok"
|
||||
|
||||
|
||||
def test_run_async_works_when_called_from_within_a_running_loop():
|
||||
"""The actual shape of a real ComfyUI node call: a synchronous function
|
||||
invoked from code that is itself already executing inside
|
||||
asyncio.run(). Must not raise "asyncio.run() cannot be called from a
|
||||
running event loop"."""
|
||||
|
||||
async def inner_coro():
|
||||
return "from inside a running loop"
|
||||
|
||||
def sync_node_function():
|
||||
# This is exactly what comfydv.ollama's chat() does: a plain sync
|
||||
# function calling _run_async(some_coroutine()).
|
||||
return _run_async(inner_coro())
|
||||
|
||||
async def outer():
|
||||
# Mirrors ComfyUI's execution.py calling `result = f(**inputs)`
|
||||
# synchronously from within its own already-running event loop.
|
||||
return sync_node_function()
|
||||
|
||||
result = asyncio.run(outer())
|
||||
assert result == "from inside a running loop"
|
||||
|
||||
|
||||
def test_run_async_propagates_exceptions_from_within_a_running_loop():
|
||||
async def failing_coro():
|
||||
raise ValueError("boom")
|
||||
|
||||
async def outer():
|
||||
return _run_async(failing_coro())
|
||||
|
||||
with pytest.raises(ValueError, match="boom"):
|
||||
asyncio.run(outer())
|
||||
|
||||
Reference in New Issue
Block a user