An application consists of a main process and an API process.
The API process launches Uvicorn threads inside it to do API related work. The main process responds to the command requested by the API process, compares the result of communicating with the outside of the application and the application status, and returns the result.
Is it possible to communicate with Queue etc. between the Uvicorn thread and the main process?
Currently, I run com_server in a separate process and communicate with com_client.
9/22/2022
Modified the com_server source code due to some behavior issues. This ensures stable operation without reaching resource limits.
com_server
import asyncio
from queue import Queue
from threading import Thread
import time
srv_addr = '127.0.0.1'
srv_port = 5638
class com_server_protcol(asyncio.Protocol):
def __init__(self, recv_queue, send_queue):
self.transport = None
self.send_queue = send_queue
self.recv_queue = recv_queue
super().__init__()
return
def connection_made(self, transport):
self.transport = transport
client_address, client_port = self.transport.get_extra_info('peername')
# ホストマシン以外の接続は拒否
if client_address != srv_addr:
self.transport.close()
def data_received(self, data):
time_out = 3.0
interval = 0.03
isLoop = True
try:
self.send_queue.put(data)
timer = time.time()
while (isLoop):
time.sleep(interval)
while (not self.recv_queue.empty()):
item = self.recv_queue.get()
if isinstance(item, str):
item = bytes(item, encoding='utf-8')
self.transport.write(item)
self.transport.close()
isLoop = False
timer = None
if timer is not None:
if time.time() - timer > time_out:
self.transport.close()
isLoop = False
except Exception as e:
print(f'com_server_protcol:{e}')
def connection_lost(self, exc):
self.transport.close()
async def com_server(srv_addr, srv_port, recv_queue, send_queue):
loop = asyncio.get_running_loop()
server = await loop.create_server(
lambda: com_server_protcol(recv_queue, send_queue),
srv_addr, srv_port)
async with server:
await server.serve_forever()
def server_run(srv_addr, srv_port, recv_queue, send_queue):
asyncio.run(com_server(srv_addr, srv_port, recv_queue, send_queue))
com_client
import asyncio
srv_addr = '127.0.0.1'
srv_port = 5638
async def client_message(message):
reader, writer = await asyncio.open_connection(srv_addr, srv_port)
writer.write(message.encode())
await writer.drain()
data = await reader.read()
writer.close()
await writer.wait_closed()
return data
def send_message(message):
recv_data = asyncio.run(client_message(message))
return recv_data
Some code for the API
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
import json
import os
import sys
current_path = os.getcwd() + "/source"
path = current_path + "/com"
sys.path.append(path)
import com_client as com # nopep8
class ApiWork(BaseModel):
category: str
router = APIRouter(
prefix=f"/work",
tags=["work"],
)
@router.post("/{work_name}")
def post_api_work(work_name: str, req_item: ApiWork):
msg_dict = {
"type": "POST",
"work": {
"name": work_name,
"category": req_item.category
}
}
recv_data = com.send_message(json.dumps(msg_dict))
if recv_data:
s_result = recv_data.decode()
if isinstance(s_result, str):
result = json.loads(s_result)
if 'status' in result:
if result['status'] == "ERROR":
raise HTTPException(status_code=400, detail=result)
if result:
return result
ret_dict = {
"status": "ERROR",
"result": {},
"error": {
"reason": "NOT RESPONSE"
}
}
raise HTTPException(status_code=400, detail=ret_dict)