2023-04-23 20:52:43 +02:00
|
|
|
import asyncio
|
2023-05-10 03:49:39 +02:00
|
|
|
import json
|
2023-04-23 20:52:43 +02:00
|
|
|
from threading import Thread
|
|
|
|
|
2023-05-10 03:49:39 +02:00
|
|
|
from websockets.server import serve
|
2023-04-23 20:52:43 +02:00
|
|
|
|
|
|
|
from extensions.api.util import build_parameters, try_start_cloudflared
|
2023-05-10 03:49:39 +02:00
|
|
|
from modules import shared
|
|
|
|
from modules.text_generation import generate_reply
|
2023-04-23 20:52:43 +02:00
|
|
|
|
|
|
|
PATH = '/api/v1/stream'
|
|
|
|
|
|
|
|
|
|
|
|
async def _handle_connection(websocket, path):
|
|
|
|
|
|
|
|
if path != PATH:
|
|
|
|
print(f'Streaming api: unknown path: {path}')
|
|
|
|
return
|
|
|
|
|
|
|
|
async for message in websocket:
|
|
|
|
message = json.loads(message)
|
|
|
|
|
|
|
|
prompt = message['prompt']
|
|
|
|
generate_params = build_parameters(message)
|
|
|
|
stopping_strings = generate_params.pop('stopping_strings')
|
2023-05-05 23:53:03 +02:00
|
|
|
generate_params['stream'] = True
|
2023-04-23 20:52:43 +02:00
|
|
|
|
|
|
|
generator = generate_reply(
|
2023-05-11 20:37:04 +02:00
|
|
|
prompt, generate_params, stopping_strings=stopping_strings, is_chat=False)
|
2023-04-23 20:52:43 +02:00
|
|
|
|
|
|
|
# As we stream, only send the new bytes.
|
2023-05-11 22:07:20 +02:00
|
|
|
skip_index = 0
|
2023-04-23 20:52:43 +02:00
|
|
|
message_num = 0
|
|
|
|
|
|
|
|
for a in generator:
|
2023-05-11 20:37:04 +02:00
|
|
|
to_send = a[skip_index:]
|
2023-04-23 20:52:43 +02:00
|
|
|
await websocket.send(json.dumps({
|
|
|
|
'event': 'text_stream',
|
|
|
|
'message_num': message_num,
|
|
|
|
'text': to_send
|
|
|
|
}))
|
|
|
|
|
2023-04-24 08:51:32 +02:00
|
|
|
await asyncio.sleep(0)
|
|
|
|
|
2023-04-23 20:52:43 +02:00
|
|
|
skip_index += len(to_send)
|
|
|
|
message_num += 1
|
|
|
|
|
|
|
|
await websocket.send(json.dumps({
|
|
|
|
'event': 'stream_end',
|
|
|
|
'message_num': message_num
|
|
|
|
}))
|
|
|
|
|
|
|
|
|
|
|
|
async def _run(host: str, port: int):
|
2023-05-03 00:03:19 +02:00
|
|
|
async with serve(_handle_connection, host, port, ping_interval=None):
|
2023-04-23 20:52:43 +02:00
|
|
|
await asyncio.Future() # run forever
|
|
|
|
|
|
|
|
|
|
|
|
def _run_server(port: int, share: bool = False):
|
|
|
|
address = '0.0.0.0' if shared.args.listen else '127.0.0.1'
|
|
|
|
|
|
|
|
def on_start(public_url: str):
|
|
|
|
public_url = public_url.replace('https://', 'wss://')
|
|
|
|
print(f'Starting streaming server at public url {public_url}{PATH}')
|
|
|
|
|
|
|
|
if share:
|
|
|
|
try:
|
|
|
|
try_start_cloudflared(port, max_attempts=3, on_start=on_start)
|
|
|
|
except Exception as e:
|
|
|
|
print(e)
|
|
|
|
else:
|
|
|
|
print(f'Starting streaming server at ws://{address}:{port}{PATH}')
|
|
|
|
|
|
|
|
asyncio.run(_run(host=address, port=port))
|
|
|
|
|
|
|
|
|
|
|
|
def start_server(port: int, share: bool = False):
|
|
|
|
Thread(target=_run_server, args=[port, share], daemon=True).start()
|