How to communicate with main process in FastAPI+Uvicorn multi-process execution environment

Viewed 58

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)

0 Answers
Related