mirror of
https://github.com/invoke-ai/InvokeAI
synced 2024-08-30 20:32:17 +00:00
357601e2d6
author Kyle Schouviller <kyle0654@hotmail.com> 1669872800 -0800 committer Kyle Schouviller <kyle0654@hotmail.com> 1676240900 -0800 Adding base node architecture Fix type annotation errors Runs and generates, but breaks in saving session Fix default model value setting. Fix deprecation warning. Fixed node api Adding markdown docs Simplifying Generate construction in apps [nodes] A few minor changes (#2510) * Pin api-related requirements * Remove confusing extra CORS origins list * Adds response models for HTTP 200 [nodes] Adding graph_execution_state to soon replace session. Adding tests with pytest. Minor typing fixes [nodes] Fix some small output query hookups [node] Fixing some additional typing issues [nodes] Move and expand graph code. Add base item storage and sqlite implementation. Update startup to match new code [nodes] Add callbacks to item storage [nodes] Adding an InvocationContext object to use for invocations to provide easier extensibility [nodes] New execution model that handles iteration [nodes] Fixing the CLI [nodes] Adding a note to the CLI [nodes] Split processing thread into separate service [node] Add error message on node processing failure Removing old files and duplicated packages Adding python-multipart
55 lines
1.6 KiB
Python
55 lines
1.6 KiB
Python
# Copyright (c) 2022 Kyle Schouviller (https://github.com/kyle0654)
|
|
|
|
import asyncio
|
|
from queue import Empty, Queue
|
|
from typing import Any
|
|
from fastapi_events.dispatcher import dispatch
|
|
from ..services.events import EventServiceBase
|
|
import threading
|
|
|
|
class FastAPIEventService(EventServiceBase):
|
|
event_handler_id: int
|
|
__queue: Queue
|
|
__stop_event: threading.Event
|
|
|
|
def __init__(self, event_handler_id: int) -> None:
|
|
self.event_handler_id = event_handler_id
|
|
self.__queue = Queue()
|
|
self.__stop_event = threading.Event()
|
|
asyncio.create_task(self.__dispatch_from_queue(stop_event = self.__stop_event))
|
|
|
|
super().__init__()
|
|
|
|
|
|
def stop(self, *args, **kwargs):
|
|
self.__stop_event.set()
|
|
self.__queue.put(None)
|
|
|
|
|
|
def dispatch(self, event_name: str, payload: Any) -> None:
|
|
self.__queue.put(dict(
|
|
event_name = event_name,
|
|
payload = payload
|
|
))
|
|
|
|
|
|
async def __dispatch_from_queue(self, stop_event: threading.Event):
|
|
"""Get events on from the queue and dispatch them, from the correct thread"""
|
|
while not stop_event.is_set():
|
|
try:
|
|
event = self.__queue.get(block = False)
|
|
if not event: # Probably stopping
|
|
continue
|
|
|
|
dispatch(
|
|
event.get('event_name'),
|
|
payload = event.get('payload'),
|
|
middleware_id = self.event_handler_id)
|
|
|
|
except Empty:
|
|
await asyncio.sleep(0.001)
|
|
pass
|
|
|
|
except asyncio.CancelledError as e:
|
|
raise e # Raise a proper error
|