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

Socket操作中SIGPIPE来源:多线程信使服务器崩溃问题排查

问题描述

我用pthreads编写了一个简单的多线程信使程序,将客户端的sockfd存入线程安全链表(链表实现取自《C++ Concurrency in Action》)。客户端处理线程会将消息转发给所有其他客户端,同时期望自动移除失效连接。但用lldb调试时发现,客户端断开连接后,服务器会因SIGPIPE(信号13)崩溃,而非预期中write返回-1,求帮忙排查原因。

服务器代码(server.cpp)

#include <pthread.h>
#include <iostream>
#include <sys/socket.h>
#include <netinet/in.h>
#include <unistd.h>
#include <cstring>
#include <stdlib.h>
#include <stdio.h>
#include <functional>
#include <signal.h>
// 假设multithread_list是来自《C++ Concurrency in Action》的线程安全链表实现
#include "multithread_list.h"

void *conn_handler(void *);
void send_msg(std::pair<uint32_t, char *> reply_pair);
multithread_list<int> adress_list;

// 自定义消息结构体(原代码未给出,补充以保证代码完整性)
struct data {
    char* name;
    uint32_t name_size;
    char* body;
    uint32_t body_size;
};

// 序列化/反序列化函数(原代码未给出,仅作声明)
data* buffer_to_data(char* buffer);
std::pair<uint32_t, char*> data_to_buffer(data* user_data);

int main(int argc, char *argv[]) {
    (void)argc;
    (void)argv;
    int sockfd;
    uint16_t portno;
    unsigned int clilen;
    struct sockaddr_in serv_addr, cli_addr;

    /* 创建socket */
    sockfd = socket(AF_INET, SOCK_STREAM, 0);

    if (sockfd < 0) {
        perror("ERROR opening socket");
        exit(1);
    }

    /* 初始化socket结构 */
    bzero((char *)&serv_addr, sizeof(serv_addr));
    portno = 5001;

    serv_addr.sin_family = AF_INET;
    serv_addr.sin_addr.s_addr = INADDR_ANY;
    serv_addr.sin_port = htons(portno);

    /* 绑定地址 */
    if (bind(sockfd, (struct sockaddr *)&serv_addr, sizeof(serv_addr)) < 0) {
        perror("ERROR on binding");
        exit(1);
    }

    /* 监听连接 */
    while (true) {
        listen(sockfd, 5);
        clilen = sizeof(cli_addr);

        /* 接受客户端连接 */
        int newsockfd = accept(sockfd, (struct sockaddr *)&cli_addr, &clilen);
        adress_list.push_front(newsockfd);
        pthread_t thread;

        if (pthread_create(&thread, nullptr, conn_handler, (void *)&newsockfd) < 0) {
            perror("ERROR creating thread.");
        }
        pthread_detach(thread);
    }
    return 0;
}

void *conn_handler(void *newsockfd_v) {
    int newsockfd = *reinterpret_cast<int *>(newsockfd_v);
    char buffer[256];
    if (newsockfd < 0) {
        perror("ERROR on accept");
        exit(1);
    }
    while (true) {
        /* 读取客户端消息 */
        bzero(buffer, 256);
        ssize_t n = read(newsockfd, buffer, 255);
        if (n < 0) {
            printf("break");
            break ;
        }
        data *user_data = buffer_to_data(buffer);
        printf("Here is the message: %s\n", user_data->body);
        printf("current newsockfd = %d\n", newsockfd);
        /* 转发消息给所有客户端 */
        send_msg(data_to_buffer(user_data));
    }
    return nullptr;
}

void send_msg(std::pair<uint32_t, char *> reply_pair) {
    adress_list.remove_if([reply_pair](int sock) {
        ssize_t n = write(sock, reply_pair.second, reply_pair.first);
        if (n < 0) {
            return true;
        }
        return false;
    });
}

客户端代码(client.cpp)

#include <pthread.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <netdb.h>
#include <unistd.h>
#include <cstring>
#include <stdlib.h>
#include <stdio.h>

#define BUFF_SIZE 256

// 自定义结构体(原代码未给出,补充以保证代码完整性)
struct data {
    char* name;
    uint32_t name_size;
    char* body;
    uint32_t body_size;
};

struct server_data {
    char* date;
    char* name;
    char* body;
};

// 辅助函数声明
void remove_enter_symbol(char *buffer, size_t buff_size);
size_t get_word_size(const char *buffer, size_t buff_size);
char* struct_to_buffer(data* user_data);
server_data* buffer_to_struct(char* buffer);

void remove_enter_symbol(char *buffer, size_t buff_size) {
    if (buffer == nullptr) exit(1);
    for (size_t i = 0; i < buff_size; i++) {
        if (buffer[i] == '\n') {
            buffer[i] = '\0';
            break;
        }
    }
}

size_t get_word_size(const char *buffer, size_t buff_size) {
    if (buffer == nullptr) exit(1);
    size_t counter = 0;
    for (size_t i = 0; i < buff_size; ++i) {
        if (buffer[i] != '\0')
            counter++;
        else
            break;
    }
    return counter;
}

void *listen_server(void *sock) {
    int sockfd = *reinterpret_cast<int *>(sock);
    char *receive_buffer = new char[256];
    while (true) {
        bzero(receive_buffer, BUFF_SIZE);
        ssize_t n = read(sockfd, receive_buffer, BUFF_SIZE - 1);
        if (n < 0) {
            perror("ERROR reading from socket");
            break;
        }
        server_data *d = buffer_to_struct(receive_buffer);
        printf("[%s]  ", d->date);
        printf("%s: ", d->name);
        printf("%s\n", d->body);
    }
    return nullptr;
}

int main(int argc, char *argv[]) {
    (void)argc;
    (void)argv;
    ssize_t sockfd, n;
    uint16_t portno;
    struct sockaddr_in serv_addr;
    struct hostent *server;

    char name_buffer[BUFF_SIZE];
    char body_buffer[BUFF_SIZE];
    data user_data;

    if (argc < 3) {
        fprintf(stderr, "usage %s hostname port\n", argv[0]);
        exit(0);
    }

    portno = (uint16_t)atoi(argv[2]);

    /* 创建socket */
    sockfd = socket(AF_INET, SOCK_STREAM, 0);

    if (sockfd < 0) {
        perror("ERROR opening socket");
        exit(1);
    }

    server = gethostbyname(argv[1]);

    if (server == NULL) {
        fprintf(stderr, "ERROR, no such host\n");
        exit(0);
    }

    bzero((char *)&serv_addr, sizeof(serv_addr));
    serv_addr.sin_family = AF_INET;
    bcopy(server->h_addr, (char *)&serv_addr.sin_addr.s_addr,
          (size_t)server->h_length);
    serv_addr.sin_port = htons(portno);

    /* 连接服务器 */
    if (connect(sockfd, (struct sockaddr *)&serv_addr, sizeof(serv_addr)) < 0) {
        perror("ERROR connecting");
        exit(1);
    }

    /* 创建监听服务器消息的线程 */
    pthread_t thread;

    if (pthread_create(&thread, nullptr, listen_server, (void *)&sockfd) < 0) {
        perror("ERROR creating thread.");
    }
    pthread_detach(thread);

    printf("Please enter the name: ");
    bzero(name_buffer, BUFF_SIZE);
    if (fgets(name_buffer, BUFF_SIZE - 1, stdin) == NULL) {
        perror("ERROR reading from stdin");
        pthread_kill(thread, -1);
        exit(1);
    }
    remove_enter_symbol(name_buffer, BUFF_SIZE);
    user_data.name = name_buffer;
    user_data.name_size = get_word_size(name_buffer, BUFF_SIZE);

    while (true) {
        printf("Please enter the message: ");
        bzero(body_buffer, BUFF_SIZE);
        if (fgets(body_buffer, BUFF_SIZE - 1, stdin) == NULL) {
            perror("ERROR reading from stdin");
            pthread_kill(thread, -1);
            exit(1);
        }
        remove_enter_symbol(body_buffer, BUFF_SIZE);
        user_data.body = body_buffer;
        user_data.body_size = get_word_size(body_buffer, BUFF_SIZE);

        char *buff_to_send = struct_to_buffer(&user_data);

        /* 发送消息到服务器 */
        n = write(sockfd, buff_to_send, BUFF_SIZE);

        if (n < 0) {
            perror("ERROR writing to socket");
            pthread_kill(thread, -1);
            exit(1);
        }
    }
    return 0;
}
问题分析与解决

核心原因:SIGPIPE信号的默认行为

当向已被客户端关闭的TCP连接写数据时,操作系统会触发SIGPIPE信号。该信号的默认处理逻辑是直接终止进程,导致程序还没等到write返回-1就崩溃。

解决方法

1. 全局忽略SIGPIPE信号

在服务器main函数开头添加以下代码,让进程忽略SIGPIPE信号:

signal(SIGPIPE, SIG_IGN);

处理后,向失效连接写数据时write会返回-1,errno被设置为EPIPE,你的remove_if逻辑就能正常移除失效的sockfd。

2. 给单个socket禁用SIGPIPE

如果不想全局忽略SIGPIPE(避免影响其他依赖该信号的逻辑),可以给每个accept得到的sockfd设置SO_NOSIGPIPE选项:

// 在accept之后添加
int opt = 1;
setsockopt(newsockfd, SOL_SOCKET, SO_NOSIGPIPE, &opt, sizeof(opt));

该选项仅对当前socket生效,向其写数据时不会触发SIGPIPE,而是让write返回-1并设置errno=EPIPE。

额外代码问题修复

除SIGPIPE外,代码还有几个潜在bug需要修正:

  • 线程参数传递错误:main中pthread_create传递的是局部变量newsockfd的指针,while循环会重复赋值该变量,新线程可能读到错误值。正确做法是把sockfd转成void*直接传递:
    // 创建线程时
    pthread_create(&thread, nullptr, conn_handler, (void*)(intptr_t)newsockfd);
    
    // 处理线程中接收参数
    void *conn_handler(void *newsockfd_v) {
        int newsockfd = (int)(intptr_t)newsockfd_v;
        // ... 其他代码
    }
    
  • 未处理read返回0的情况:客户端正常关闭连接时read会返回0,当前代码会进入无限循环反复读取0,需添加n==0的判断并退出循环。
  • 内存泄漏:data_to_buffer和struct_to_buffer分配的内存需手动释放,避免内存泄漏。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.17 21:46:06