Compare commits
4
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d735d8ae9c | ||
|
|
349aee5420 | ||
|
|
0ae70835c4 | ||
|
|
99cda529a9 |
+40
-9
@@ -64,14 +64,43 @@ async def ensure_client_session():
|
||||
|
||||
async def cleanup():
|
||||
global client_session
|
||||
if client_session:
|
||||
await client_session.close()
|
||||
try:
|
||||
# Cancel the monitor task if it exists
|
||||
if hasattr(upload_queue, "monitor_task") and upload_queue.monitor_task:
|
||||
if not upload_queue.monitor_task.done():
|
||||
upload_queue.monitor_task.cancel()
|
||||
try:
|
||||
await upload_queue.monitor_task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
except RuntimeError as e:
|
||||
# Handle case where event loop is closed
|
||||
if "Event loop is closed" in str(e):
|
||||
logger.info("Event loop closed during cleanup")
|
||||
else:
|
||||
raise
|
||||
logger.info("Upload queue monitor task cancelled")
|
||||
|
||||
# Clean up the client session
|
||||
if client_session:
|
||||
await client_session.close()
|
||||
|
||||
# Clean up the upload queue
|
||||
# await upload_queue.cleanup()
|
||||
except Exception as e:
|
||||
logger.error(f"Error during cleanup: {str(e)}")
|
||||
|
||||
|
||||
def exit_handler():
|
||||
print("Exiting the application. Initiating cleanup...")
|
||||
loop = asyncio.get_event_loop()
|
||||
loop.run_until_complete(cleanup())
|
||||
try:
|
||||
loop = asyncio.get_event_loop()
|
||||
if loop.is_closed():
|
||||
logger.warning("Event loop is already closed during exit")
|
||||
return
|
||||
loop.run_until_complete(cleanup())
|
||||
except Exception as e:
|
||||
print(f"Error during exit cleanup: {str(e)}")
|
||||
|
||||
|
||||
atexit.register(exit_handler)
|
||||
@@ -139,7 +168,9 @@ async def async_request_with_retry(
|
||||
logger.error(f"Error response body: {error_body}")
|
||||
|
||||
if attempt == max_retries - 1:
|
||||
logger.error(f"Request failed after {max_retries} attempts: {e}")
|
||||
logger.error(
|
||||
f"Request {method} : {url} failed after {max_retries} attempts: {e}"
|
||||
)
|
||||
raise
|
||||
|
||||
await asyncio.sleep(retry_delay)
|
||||
@@ -147,7 +178,7 @@ async def async_request_with_retry(
|
||||
|
||||
total_time = time.time() - start_time
|
||||
raise Exception(
|
||||
f"Request failed after {max_retries} attempts and {total_time:.2f} seconds"
|
||||
f"Request {method} : {url} failed after {max_retries} attempts and {total_time:.2f} seconds"
|
||||
)
|
||||
|
||||
|
||||
@@ -1310,7 +1341,7 @@ async def send_json_override(self, event, data, sid=None):
|
||||
if prompt_id in comfy_message_queues:
|
||||
comfy_message_queues[prompt_id].put_nowait({"event": event, "data": data})
|
||||
|
||||
asyncio.create_task(update_run_ws_event(prompt_id, event, data))
|
||||
# asyncio.create_task(update_run_ws_event(prompt_id, event, data))
|
||||
|
||||
if event == "execution_start":
|
||||
if prompt_id in prompt_metadata:
|
||||
@@ -2167,8 +2198,8 @@ async def initialize_upload_queue(app=None):
|
||||
upload_queue.max_concurrent,
|
||||
)
|
||||
|
||||
# Start the queue monitoring task in the same event loop
|
||||
asyncio.create_task(monitor_upload_queue())
|
||||
# Store the monitor task reference
|
||||
upload_queue.monitor_task = asyncio.create_task(monitor_upload_queue())
|
||||
|
||||
|
||||
# Get the server's event loop and initialize there
|
||||
|
||||
Reference in New Issue
Block a user