You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Asyncio StreamWriter wait_closed阻塞代码导致程序挂起的问题求助

Asyncio StreamWriter wait_closed阻塞代码导致程序挂起的问题求助

我写了一个非常简陋的服务器demo,想复现我遇到的问题:当我尝试断开客户端连接时,程序会永久挂起,我完全搞不懂原因。一开始我怀疑是read方法的问题,但就算加上了超时实现,问题还是存在。

import asyncio
from asyncio import StreamReader, StreamWriter, wait_for, start_server
import logging
import socket

logging.basicConfig(
    level=logging.DEBUG,
    format="%(asctime)s - %(levelname)s - %(message)s",
)


class Client:
    def __init__(self, reader: StreamReader, writer: StreamWriter) -> None:
        self.reader = reader
        self.writer = writer
        self.address: str = ":".join(map(str, writer.get_extra_info("peername")))
        self._configure_tcp()
        logging.info(f"New connection: {self.address}")

    async def write(self, data: bytes) -> None:
        logging.debug(f"Sending to {self.address}: {data}")
        self.writer.write(data)
        await self.writer.drain()

    async def read(self) -> bytes:
        try:
            if data := await wait_for(self.reader.read(1024), timeout=1):
                logging.debug(f"Received from {self.address}: {data}")
                self.last_msg = data
                return data
            else:
                return b""
        except TimeoutError:
            return b""

    async def close(self) -> None:
        if not self.writer.is_closing():
            self.writer.close()
            await self.writer.wait_closed()
            logging.info(f"Connection closed: {self.address}")

    def _configure_tcp(self) -> None:
        _socket = self.writer.get_extra_info("socket")
        _socket.setsockopt(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)
        _socket.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, 10)
        _socket.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, 10)
        _socket.setsockopt(socket.IPPROTO_TCP, socket.TCP_KEEPCNT, 5)


class ServerProtocol:
    def __init__(self, name: str, port: int):
        self.name: str = name
        self.port: int = port
        self.clients: list[Client] = []
        self.server = None

    async def start_server(self) -> None:
        self.server = await start_server(self.connection, "0.0.0.0", self.port)
        logging.info(f"Serving at 0.0.0.0:{self.port}")

    async def stop_server(self) -> None:
        if self.server:
            logging.info("Stoping server...")
            self.server.close()
            await self.server.wait_closed()
            logging.info("Server finalized")

        for client in self.clients:
            await self.disconnect(client)
        logging.info("All connections closed")

    async def connection(self, reader: StreamReader, writer: StreamWriter) -> None:
        client = Client(reader, writer)
        self.clients.append(client)
        try:
            while True:
                data = await client.read()
                if data:
                    await self._process(client, data)
        except TimeoutError:
            logging.info(f"Timeout: {client.address}")
        except Exception as e:
            logging.error(f"Exception with {client.address}: {e}")
        finally:
            await self.disconnect(client)

    async def disconnect(self, client: Client) -> None:
        logging.info(f"Disconnecting {client.address}")
        if client in self.clients:
            self.clients.remove(client)
            await client.close()
            logging.info(f"{client.address} disconnected.")

    async def _process(self, client: Client, data: bytes) -> None:
        await client.write(data)


async def main():
    server = ServerProtocol(name="TestServer", port=8888)

    try:
        await server.start_server()
        await asyncio.sleep(30)
        await server.disconnect(server.clients[0])
        await server.stop_server()
    except KeyboardInterrupt:
        await server.stop_server()


if __name__ == "__main__":
    asyncio.run(main())

这个demo的逻辑很简单:监听8888端口,等待30秒后,获取第一个连接的客户端并尝试断开它,之后再关闭服务器。另外,就算是我用nc连接服务器后主动关闭客户端(比如关掉nc窗口),同样会出现程序挂起的问题。

有没有大佬能帮忙分析下问题出在哪?

备注:内容来源于stack exchange,提问作者IamRichter

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.04.13 19:09:33