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

Windows平台C++多客户端全双工命名管道通信问题排查

Windows平台C++服务端+C#客户端全双工命名管道通信方案实现

一、单管道异步全双工方案(优先推荐)

核心思路

使用Windows命名管道的异步IO模型(重叠IO),服务端为每个客户端连接创建独立的IO处理线程,同时异步处理读、写操作,彻底解决同步IO导致的"必须先收再发"和缓存延迟问题,真正实现全双工通信。

C++服务端代码(支持多客户端+异步全双工)

#include <windows.h>
#include <iostream>
#include <vector>
#include <thread>
#include <mutex>

#define PIPE_NAME L"\\\\.\\pipe\\AsyncDuplexPipe"
#define BUFFER_SIZE 4096

// 每个客户端的上下文数据
typedef struct {
    HANDLE hPipe;
    OVERLAPPED overlappedRead;
    OVERLAPPED overlappedWrite;
    char readBuffer[BUFFER_SIZE];
    char writeBuffer[BUFFER_SIZE];
    BOOL isWriting;
    BOOL isReading;
} ClientContext;

std::vector<ClientContext*> clients;
std::mutex clientsMutex;

// 处理客户端连接的线程函数
void ClientHandler(ClientContext* pCtx) {
    DWORD bytesTransferred;
    BOOL result;

    // 初始化异步读
    pCtx->isReading = TRUE;
    ZeroMemory(&pCtx->overlappedRead, sizeof(OVERLAPPED));
    pCtx->overlappedRead.hEvent = CreateEvent(NULL, TRUE, FALSE, NULL);
    result = ReadFile(pCtx->hPipe, pCtx->readBuffer, BUFFER_SIZE, NULL, &pCtx->overlappedRead);

    if (!result && GetLastError() != ERROR_IO_PENDING) {
        std::cerr << "ReadFile failed: " << GetLastError() << std::endl;
        goto Cleanup;
    }

    // 循环处理IO操作
    while (TRUE) {
        // 等待读或写完成
        HANDLE events[] = { pCtx->overlappedRead.hEvent, pCtx->overlappedWrite.hEvent };
        DWORD waitResult = WaitForMultipleObjects(2, events, FALSE, INFINITE);

        switch (waitResult) {
            case WAIT_OBJECT_0: // 读操作完成
                ResetEvent(pCtx->overlappedRead.hEvent);
                if (!GetOverlappedResult(pCtx->hPipe, &pCtx->overlappedRead, &bytesTransferred, FALSE)) {
                    std::cerr << "Read failed: " << GetLastError() << std::endl;
                    goto Cleanup;
                }
                if (bytesTransferred == 0) { // 客户端断开
                    std::cout << "Client disconnected" << std::endl;
                    goto Cleanup;
                }

                // 处理收到的消息(示例:回显+主动发送响应)
                pCtx->readBuffer[bytesTransferred] = '\0';
                std::cout << "Received from client: " << pCtx->readBuffer << std::endl;

                // 主动发送消息(无需等待读操作完成)
                if (!pCtx->isWriting) {
                    snprintf(pCtx->writeBuffer, BUFFER_SIZE, "Server response: %s", pCtx->readBuffer);
                    pCtx->isWriting = TRUE;
                    ZeroMemory(&pCtx->overlappedWrite, sizeof(OVERLAPPED));
                    pCtx->overlappedWrite.hEvent = CreateEvent(NULL, TRUE, FALSE, NULL);
                    result = WriteFile(pCtx->hPipe, pCtx->writeBuffer, strlen(pCtx->writeBuffer), NULL, &pCtx->overlappedWrite);
                    if (!result && GetLastError() != ERROR_IO_PENDING) {
                        std::cerr << "WriteFile failed: " << GetLastError() << std::endl;
                        goto Cleanup;
                    }
                }

                // 重新发起异步读
                ZeroMemory(pCtx->readBuffer, BUFFER_SIZE);
                result = ReadFile(pCtx->hPipe, pCtx->readBuffer, BUFFER_SIZE, NULL, &pCtx->overlappedRead);
                if (!result && GetLastError() != ERROR_IO_PENDING) {
                    std::cerr << "ReadFile failed: " << GetLastError() << std::endl;
                    goto Cleanup;
                }
                break;

            case WAIT_OBJECT_0 + 1: // 写操作完成
                ResetEvent(pCtx->overlappedWrite.hEvent);
                if (!GetOverlappedResult(pCtx->hPipe, &pCtx->overlappedWrite, &bytesTransferred, FALSE)) {
                    std::cerr << "Write failed: " << GetLastError() << std::endl;
                    goto Cleanup;
                }
                pCtx->isWriting = FALSE;
                CloseHandle(pCtx->overlappedWrite.hEvent);
                break;

            default:
                std::cerr << "Wait failed: " << GetLastError() << std::endl;
                goto Cleanup;
        }
    }

Cleanup:
    // 清理资源
    CloseHandle(pCtx->overlappedRead.hEvent);
    if (pCtx->isWriting) CloseHandle(pCtx->overlappedWrite.hEvent);
    CloseHandle(pCtx->hPipe);

    // 从客户端列表移除
    std::lock_guard<std::mutex> lock(clientsMutex);
    auto it = std::find(clients.begin(), clients.end(), pCtx);
    if (it != clients.end()) clients.erase(it);
    delete pCtx;
}

// 监听客户端连接的主线程
int main() {
    while (TRUE) {
        HANDLE hPipe = CreateNamedPipe(
            PIPE_NAME,
            PIPE_ACCESS_DUPLEX | FILE_FLAG_OVERLAPPED, // 双工+重叠IO
            PIPE_TYPE_MESSAGE | PIPE_READMODE_MESSAGE | PIPE_WAIT,
            PIPE_UNLIMITED_INSTANCES, // 支持多客户端
            BUFFER_SIZE,
            BUFFER_SIZE,
            0,
            NULL
        );

        if (hPipe == INVALID_HANDLE_VALUE) {
            std::cerr << "CreateNamedPipe failed: " << GetLastError() << std::endl;
            return 1;
        }

        // 等待客户端连接
        BOOL connected = ConnectNamedPipe(hPipe, NULL);
        if (!connected && GetLastError() != ERROR_PIPE_CONNECTED) {
            std::cerr << "ConnectNamedPipe failed: " << GetLastError() << std::endl;
            CloseHandle(hPipe);
            continue;
        }

        std::cout << "New client connected" << std::endl;

        // 创建客户端上下文
        ClientContext* pCtx = new ClientContext();
        pCtx->hPipe = hPipe;
        pCtx->isWriting = FALSE;

        // 添加到客户端列表
        std::lock_guard<std::mutex> lock(clientsMutex);
        clients.push_back(pCtx);

        // 启动线程处理该客户端
        std::thread handlerThread(ClientHandler, pCtx);
        handlerThread.detach();
    }

    return 0;
}

C#客户端代码(异步全双工)

using System;
using System.IO.Pipes;
using System.Text;
using System.Threading.Tasks;

class AsyncDuplexPipeClient
{
    static async Task Main(string[] args)
    {
        using (var pipeClient = new NamedPipeClientStream(".", "AsyncDuplexPipe", PipeDirection.InOut, PipeOptions.Asynchronous))
        {
            await pipeClient.ConnectAsync();
            Console.WriteLine("Connected to server");

            // 同时启动读、写任务,实现全双工
            var readTask = ReadFromPipeAsync(pipeClient);
            var writeTask = WriteToPipeAsync(pipeClient);

            await Task.WhenAll(readTask, writeTask);
        }
    }

    static async Task ReadFromPipeAsync(NamedPipeClientStream pipe)
    {
        byte[] buffer = new byte[4096];
        while (pipe.IsConnected)
        {
            try
            {
                int bytesRead = await pipe.ReadAsync(buffer, 0, buffer.Length);
                if (bytesRead == 0) break;
                string message = Encoding.ASCII.GetString(buffer, 0, bytesRead);
                Console.WriteLine($"Received from server: {message}");
            }
            catch (Exception ex)
            {
                Console.WriteLine($"Read error: {ex.Message}");
                break;
            }
        }
    }

    static async Task WriteToPipeAsync(NamedPipeClientStream pipe)
    {
        while (pipe.IsConnected)
        {
            string input = Console.ReadLine();
            if (string.IsNullOrEmpty(input)) break;

            byte[] buffer = Encoding.ASCII.GetBytes(input);
            try
            {
                await pipe.WriteAsync(buffer, 0, buffer.Length);
                await pipe.FlushAsync();
            }
            catch (Exception ex)
            {
                Console.WriteLine($"Write error: {ex.Message}");
                break;
            }
        }
    }
}

二、双管道方案(解决周期性断开问题)

核心思路

为每个客户端创建两个独立管道:一个用于客户端→服务端的读管道,一个用于服务端→客户端的写管道。之前的周期性断开问题通常是因为管道句柄未正确维护、未处理异常断开场景,以下代码添加了重连逻辑和完善的句柄生命周期管理。

C++服务端代码(双管道+多客户端)

#include <windows.h>
#include <iostream>
#include <vector>
#include <thread>
#include <mutex>
#include <string>

#define READ_PIPE_PREFIX L"\\\\.\\pipe\\ClientToServer_"
#define WRITE_PIPE_PREFIX L"\\\\.\\pipe\\ServerToClient_"
#define BUFFER_SIZE 4096

typedef struct {
    std::string clientId;
    HANDLE hReadPipe;
    HANDLE hWritePipe;
    BOOL isActive;
} DualPipeClient;

std::vector<DualPipeClient*> clients;
std::mutex clientsMutex;

void ReadThread(DualPipeClient* pClient) {
    char buffer[BUFFER_SIZE];
    DWORD bytesRead;

    while (pClient->isActive) {
        if (!ReadFile(pClient->hReadPipe, buffer, BUFFER_SIZE, &bytesRead, NULL)) {
            if (GetLastError() == ERROR_BROKEN_PIPE) {
                std::cout << "Client " << pClient->clientId << " read pipe disconnected" << std::endl;
                break;
            }
            std::cerr << "Read failed: " << GetLastError() << std::endl;
            break;
        }
        if (bytesRead == 0) break;

        buffer[bytesRead] = '\0';
        std::cout << "Client " << pClient->clientId << " sent: " << buffer << std::endl;

        // 主动通过写管道发送响应
        std::string response = "Server ack: " + std::string(buffer);
        DWORD bytesWritten;
        WriteFile(pClient->hWritePipe, response.c_str(), response.length(), &bytesWritten, NULL);
    }

    pClient->isActive = FALSE;
}

void WriteThread(DualPipeClient* pClient) {
    // 模拟周期性主动发送心跳消息
    while (pClient->isActive) {
        std::string message = "Server heartbeat to " + pClient->clientId;
        DWORD bytesWritten;
        if (!WriteFile(pClient->hWritePipe, message.c_str(), message.length(), &bytesWritten, NULL)) {
            if (GetLastError() == ERROR_BROKEN_PIPE) {
                std::cout << "Client " << pClient->clientId << " write pipe disconnected" << std::endl;
                break;
            }
            std::cerr << "Write failed: " << GetLastError() << std::endl;
            break;
        }
        Sleep(5000);
    }

    pClient->isActive = FALSE;
}

int main() {
    int clientCounter = 0;
    while (TRUE) {
        // 生成客户端唯一ID
        std::string clientId = "Client_" + std::to_string(++clientCounter);
        std::wstring readPipeName = READ_PIPE_PREFIX + std::wstring(clientId.begin(), clientId.end());
        std::wstring writePipeName = WRITE_PIPE_PREFIX + std::wstring(clientId.begin(), clientId.end());

        // 创建读管道(客户端→服务端)
        HANDLE hReadPipe = CreateNamedPipe(
            readPipeName.c_str(),
            PIPE_ACCESS_INBOUND,
            PIPE_TYPE_MESSAGE | PIPE_READMODE_MESSAGE | PIPE_WAIT,
            PIPE_UNLIMITED_INSTANCES,
            BUFFER_SIZE,
            BUFFER_SIZE,
            0,
            NULL
        );
        if (hReadPipe == INVALID_HANDLE_VALUE) {
            std::cerr << "Create read pipe failed: " << GetLastError() << std::endl;
            continue;
        }

        // 创建写管道(服务端→客户端)
        HANDLE hWritePipe = CreateNamedPipe(
            writePipeName.c_str(),
            PIPE_ACCESS_OUTBOUND,
            PIPE_TYPE_MESSAGE | PIPE_WAIT,
            PIPE_UNLIMITED_INSTANCES,
            BUFFER_SIZE,
            BUFFER_SIZE,
            0,
            NULL
        );
        if (hWritePipe == INVALID_HANDLE_VALUE) {
            std::cerr << "Create write pipe failed: " << GetLastError() << std::endl;
            CloseHandle(hReadPipe);
            continue;
        }

        // 等待客户端连接两个管道
        std::cout << "Waiting for client " << clientId << " to connect..." << std::endl;
        BOOL readConnected = ConnectNamedPipe(hReadPipe, NULL);
        if (!readConnected && GetLastError() != ERROR_PIPE_CONNECTED) {
            std::cerr << "Connect read pipe failed: " << GetLastError() << std::endl;
            CloseHandle(hReadPipe);
            CloseHandle(hWritePipe);
            continue;
        }

        BOOL writeConnected = ConnectNamedPipe(hWritePipe, NULL);
        if (!writeConnected && GetLastError() != ERROR_PIPE_CONNECTED) {
            std::cerr << "Connect write pipe failed: " << GetLastError() << std::endl;
            CloseHandle(hReadPipe);
            CloseHandle(hWritePipe);
            continue;
        }

        std::cout << "Client " << clientId << " connected" << std::endl;

        // 保存客户端信息
        DualPipeClient* pClient = new DualPipeClient();
        pClient->clientId = clientId;
        pClient->hReadPipe = hReadPipe;
        pClient->hWritePipe = hWritePipe;
        pClient->isActive = TRUE;

        std::lock_guard<std::mutex> lock(clientsMutex);
        clients.push_back(pClient);

        // 启动读和写线程
        std::thread readThread(ReadThread, pClient);
        std::thread writeThread(WriteThread, pClient);
        readThread.detach();
        writeThread.detach();
    }

    return 0;
}

C#客户端代码(双管道+重连逻辑)

using System;
using System.IO.Pipes;
using System.Text;
using System.Threading.Tasks;

class DualPipeClient
{
    static string _clientId = "Client_" + Guid.NewGuid().ToString().Substring(0, 8);
    static NamedPipeClientStream _readPipe;
    static NamedPipeClientStream _writePipe;
    static bool _isRunning = true;

    static async Task Main(string[] args)
    {
        await ConnectPipesAsync();

        var readTask = ReadFromPipeAsync();
        var writeTask = WriteToPipeAsync();
        var heartbeatTask = KeepAliveAsync();

        await Task.WhenAll(readTask, writeTask, heartbeatTask);
    }

    static async Task ConnectPipesAsync()
    {
        while (_isRunning)
        {
            try
            {
                _readPipe = new NamedPipeClientStream(".", $"ServerToClient_{_clientId}", PipeDirection.In, PipeOptions.None);
                _writePipe = new NamedPipeClientStream(".", $"ClientToServer_{_clientId}", PipeDirection.Out, PipeOptions.None);

                await Task.WhenAll(_readPipe.ConnectAsync(1000), _writePipe.ConnectAsync(1000));
                Console.WriteLine("Connected to server dual pipes");
                return;
            }
            catch (TimeoutException)
            {
                Console.WriteLine("Waiting for server pipes...");
                await Task.Delay(1000);
            }
            catch (Exception ex)
            {
                Console.WriteLine($"Connection error: {ex.Message}");
                await Task.Delay(2000);
            }
        }
    }

    static async Task ReadFromPipeAsync()
    {
        byte[] buffer = new byte[4096];
        while (_isRunning)
        {
            try
            {
                if (!_readPipe.IsConnected)
                {
                    Console.WriteLine("Read pipe disconnected, reconnecting...");
                    await ConnectPipesAsync();
                }

                int bytesRead = await _readPipe.ReadAsync(buffer, 0, buffer.Length);
                if (bytesRead == 0) throw new InvalidOperationException("Pipe closed");

                string message = Encoding.ASCII.GetString(buffer, 0, bytesRead);
                Console.WriteLine($"Received from server
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 00:09:39