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
相关产品推荐
相关产品推荐

