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

C++20协程+Asio单线程模型中forkpty进程stdout/stderr读取阻塞问题

为何Asio协程中async_read_some阻塞而read正常工作?

我是C++协程和Asio的新手,正在开发一款本地进程管理工具:客户端负责发送程序启动命令,后端守护进程采用单线程多协程模型管理进程。守护进程需要与子进程交互获取日志,通过IPC传输给客户端。

我的问题是:为何守护进程在代码1的async_read_some处阻塞,无法触发对stdout的读取,而代码2的read操作却能正常工作?

简化演示代码

客户端代码

#include <iostream>
#include <pwd.h>
#include <sys/socket.h>
#include <sys/un.h>
#include <thread>
#include <unistd.h>

void receive_output(int sockfd)
{
    char buffer[1024];
    while (true)
    {
        ssize_t n = recv(sockfd, buffer, sizeof(buffer) - 1, 0);
        if (n <= 0)
            break;
        buffer[n] = '\0';
        std::cout << buffer << std::flush;
    }
}

void send_input(int sockfd)
{
    std::string input;
    while (std::getline(std::cin, input))
    {
        input += '\n';
        send(sockfd, input.c_str(), input.size(), 0);
    }
}

int main()
{
    int sockfd = socket(AF_UNIX, SOCK_STREAM, 0);
    struct sockaddr_un addr;
    addr.sun_family = AF_UNIX;
    strcpy(addr.sun_path, "/tmp/socket");

    connect(sockfd, (struct sockaddr *)&addr, sizeof(addr));

    std::string command = "test";
    send(sockfd, command.c_str(), command.size(), 0);

    std::thread output_thread(receive_output, sockfd);
    std::thread input_thread(send_input, sockfd);

    output_thread.join();
    input_thread.join();

    close(sockfd);
    return 0;
}

服务端代码

#include <boost/asio.hpp>
#include <boost/asio/awaitable.hpp>
#include <boost/asio/co_spawn.hpp>
#include <boost/asio/posix/stream_descriptor.hpp>
#include <boost/asio/use_awaitable.hpp>
#include <cstdlib>
#include <fcntl.h>
#include <iostream>
#include <pty.h>
#include <pwd.h>
#include <sys/socket.h>
#include <sys/un.h>
#include <sys/wait.h>

using boost::asio::awaitable;
using boost::asio::use_awaitable;
namespace asio = boost::asio;

boost::asio::io_context io_context;

void set_nonblocking(int fd)
{
    int flags = fcntl(fd, F_GETFL, 0);
    if (flags == -1)
    {
        perror("fcntl(F_GETFL)");
        return;
    }
    flags |= O_NONBLOCK;
    if (fcntl(fd, F_SETFL, flags) == -1)
    {
        perror("fcntl(F_SETFL)");
    }
}

awaitable<std::string> receive_command(asio::posix::stream_descriptor &client_stream)
{
    char buffer[256];
    std::size_t len = co_await client_stream.async_read_some(boost::asio::buffer(buffer), use_awaitable);
    buffer[len] = '\0';
    co_return std::string(buffer);
}

awaitable<void> run_process_in_pty(const std::string &command, asio::posix::stream_descriptor &client_stream,
                                   int client_fd)
{
    int master_fd;
    pid_t pid = forkpty(&master_fd, NULL, NULL, NULL);
    //   set_nonblocking(master_fd);
    if (pid == 0)
    {
        // Child process: Set user and change directory
        setvbuf(stdout, nullptr, _IONBF, 0);

        // execute echo just for test
        // execl(command.c_str(), "", NULL);
        execl("/bin/echo", "echo", "Testing PTY output!", NULL);

        std::cerr << "Failed to execute command: " << strerror(errno) << std::endl;
        exit(EXIT_FAILURE);
    }
    else
    {
        // Parent process: Handle I/O between client and pty
        asio::posix::stream_descriptor pty_stream(io_context, master_fd);

        // handle output
        co_spawn(
            io_context,
            [&client_stream, &pty_stream]() -> awaitable<void> {
                char buffer[1024];
                try
                {
                    while (true)
                    {
                        // code 1
                        // std::size_t n = co_await pty_stream.async_read_some(boost::asio::buffer(buffer), use_awaitable);

                        // code 2
                        ssize_t n = read(pty_stream.native_handle(), buffer, sizeof(buffer));

                        std::cout << buffer << std::endl;
                        if (n == 0)
                            break; // End of stream

                        //   co_await asio::async_write(
                        //   client_stream, boost::asio::buffer(buffer, n),
                        //   use_awaitable);
                    }
                }
                catch (std::exception &e)
                {
                    std::cerr << strerror(errno) << std::endl;
                }
                co_return;
            },
            asio::detached);

        // handle input

        waitpid(pid, nullptr, 0);

        co_return;
    }
}

awaitable<void> accept_client(asio::posix::stream_descriptor client_stream, int client_fd)
{
    std::string command = co_await receive_command(client_stream);
    std::cout << "Received command: " << command << std::endl;

    co_await run_process_in_pty(command, client_stream, client_fd);
    client_stream.close();
}

awaitable<int> async_accept(int server_fd)
{
    int client_fd;
    co_await asio::post(io_context, use_awaitable);
    client_fd = accept(server_fd, nullptr, nullptr);
    if (client_fd < 0)
    {
        std::cerr << "Failed to accept client connection" << std::endl;
        co_return - 1;
    }
    std::cout << "clientfd, " << client_fd << std::endl;
    set_nonblocking(client_fd);
    co_return client_fd;
}

awaitable<void> daemon_loop()
{
    const char *socket_path = "/tmp/socket";

    if (access(socket_path, F_OK) == 0)
    {
        if (unlink(socket_path) != 0)
        {
            std::cerr << "Failed to remove existing socket file" << std::endl;
            exit(EXIT_FAILURE);
        }
    }

    int server_fd = socket(AF_UNIX, SOCK_STREAM, 0);
    if (server_fd < 0)
    {
        std::cerr << "Failed to create socket" << std::endl;
        exit(EXIT_FAILURE);
    }
    //   set_nonblocking(server_fd);

    struct sockaddr_un addr;
    memset(&addr, 0, sizeof(addr));
    addr.sun_family = AF_UNIX;
    strcpy(addr.sun_path, socket_path);
    if (bind(server_fd, (struct sockaddr *)&addr, sizeof(addr)) < 0)
    {
        std::cerr << "Failed to bind socket" << std::endl;
        exit(EXIT_FAILURE);
    }

    if (listen(server_fd, 5) < 0)
    {
        std::cerr << "Failed to listen on socket" << std::endl;
        exit(EXIT_FAILURE);
    }

    asio::posix::stream_descriptor server_stream(io_context, server_fd);

    while (true)
    {
        int client_fd = co_await async_accept(server_fd);
        asio::posix::stream_descriptor client_stream(io_context, client_fd);
        std::cout << "Accepted connection" << std::endl;
        co_spawn(io_context, accept_client(std::move(client_stream), client_fd), asio::detached);
    }
}

int main()
{
    std::cout << "start" << std::endl;
    co_spawn(io_context, daemon_loop(), asio::detached);
    io_context.run();
    return 0;
}

运行结果差异

  • 代码1运行结果:代码1运行结果
  • 代码2运行结果:代码2运行结果

问题原因及修复方案

核心原因

  1. pty主设备未设置非阻塞模式
    Asio的异步操作要求文件描述符必须处于非阻塞模式。你注释掉了set_nonblocking(master_fd),导致pty主设备默认是阻塞模式:

    • 代码2的同步read:在阻塞模式下会一直等待数据,因此能正常读取子进程输出。
    • 代码1的async_read_some:Asio内部依赖非阻塞系统调用,fd为阻塞模式时会直接卡住整个io_context事件循环,协程挂起后无法被唤醒。
  2. waitpid阻塞事件循环
    父进程中直接调用waitpid(pid, nullptr, 0)是阻塞调用,会卡住当前协程的执行流,导致io_context无法处理pty_stream的异步读取事件,最终async_read_some的回调永远无法触发。

修复方案

  • 启用pty非阻塞模式:取消注释set_nonblocking(master_fd),确保Asio异步操作能正常工作。
  • 异步等待子进程退出:替换阻塞的waitpid,比如在单独的协程中使用非阻塞的waitpid(pid, nullptr, WNOHANG),配合asio::post定期检查子进程状态,避免阻塞事件循环。

内容的提问来源于stack exchange,提问作者Jarviswu

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 11:17:34