Compare commits

...
Author SHA1 Message Date
bennykok d735d8ae9c fix 2025-03-28 10:15:09 +01:00
bennykok 349aee5420 fix: upload queue cleanup 2025-03-28 09:55:02 +01:00
BennyKok 0ae70835c4 remove sending ws event, since its not used. 2025-03-27 00:04:34 +08:00
BennyKok 99cda529a9 chore: clearer log when request time out 2025-03-26 23:55:47 +08:00
+40 -9
View File
@@ -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