Skip to content

Base::PackBaseMsg应该提供返回值告知调用者是否成功打包消息 #280

Description

@dyx2025

Base::PackBaseMsg应该提供返回值告知调用者是否成功打包消息

Base::PackBaseMsg没有返回值,无法告知调用者是否成功打包消息。Base::PackBaseMsg在Header序列化失败时要么程序崩溃(没有NDEBUG宏定义,目前phxpaxos的情况),要么继续运行(有NDEBUG宏定义)。

序列化PaxosMsg和CheckpointMsg后,紧接着序列化Header。PaxosMsg和CheckpointMsg序列化失败,则返回错误码。而Header序列化失败,没有返回错误码,处理逻辑不一致。

Header序列化失败时要么程序崩溃,要么继续运行,都不是好的做法。程序崩溃会拖慢前端模块和phxpaxos内部模块处理请求的速度。继续运行会把没有正确打包的消息传播到其他paxos节点。

基于以上理由,Base::PackBaseMsg应该提供返回值告知调用者是否成功打包消息。

原代码路径:
src/algorithm/base.h

class Base
{
    // ......    
    void PackBaseMsg(const std::string & sBodyBuffer, const int iCmd, std::string & sBuffer);
    // ......  
};

src/algorithm/base.cpp

int Base :: PackMsg(const PaxosMsg & oPaxosMsg, std::string & sBuffer)
{
     // PaxosMsg序列化失败时,返回错误码
    std::string sBodyBuffer;
    bool bSucc = oPaxosMsg.SerializeToString(&sBodyBuffer);
    if (!bSucc)
    {   
        PLGErr("PaxosMsg.SerializeToString fail, skip this msg");
        return -1;
    }
    
    int iCmd = MsgCmd_PaxosMsg;
    PackBaseMsg(sBodyBuffer, iCmd, sBuffer);
    
    return 0;
}

int Base :: PackCheckpointMsg(const CheckpointMsg & oCheckpointMsg, std::string & sBuffer)
{
    // CheckpointMsg序列化失败时,返回错误码
    std::string sBodyBuffer;
    bool bSucc = oCheckpointMsg.SerializeToString(&sBodyBuffer);
    if (!bSucc)
    {
        PLGErr("CheckpointMsg.SerializeToString fail, skip this msg");
        return -1;
    }

    int iCmd = MsgCmd_CheckpointMsg;
    PackBaseMsg(sBodyBuffer, iCmd, sBuffer);

    return 0;
}

void Base :: PackBaseMsg(const std::string & sBodyBuffer, const int iCmd, std::string & sBuffer)
{
    char sGroupIdx[GROUPIDXLEN] = {0};
    int iGroupIdx = m_poConfig->GetMyGroupIdx();
    memcpy(sGroupIdx, &iGroupIdx, sizeof(sGroupIdx));

    Header oHeader;
    oHeader.set_gid(m_poConfig->GetGid());
    oHeader.set_rid(0);
    oHeader.set_cmdid(iCmd);
    oHeader.set_version(1);

    // Header序列化失败时,要么程序崩溃,要么继续运行
    std::string sHeaderBuffer;
    bool bSucc = oHeader.SerializeToString(&sHeaderBuffer);
    if (!bSucc)
    {
        PLGErr("Header.SerializeToString fail, skip this msg");
        assert(bSucc == true);
    }

    char sHeaderLen[HEADLEN_LEN] = {0};
    uint16_t iHeaderLen = (uint16_t)sHeaderBuffer.size();
    memcpy(sHeaderLen, &iHeaderLen, sizeof(sHeaderLen));

    sBuffer = string(sGroupIdx, sizeof(sGroupIdx)) + string(sHeaderLen, sizeof(sHeaderLen)) + sHeaderBuffer + sBodyBuffer;

    //check sum
    uint32_t iBufferChecksum = crc32(0, (const uint8_t *)sBuffer.data(), sBuffer.size(), NET_CRC32SKIP);
    char sBufferChecksum[CHECKSUM_LEN] = {0};
    memcpy(sBufferChecksum, &iBufferChecksum, sizeof(sBufferChecksum));

    sBuffer += string(sBufferChecksum, sizeof(sBufferChecksum));
}

修改后的代码:
src/algorithm/base.h

class Base
{
    // ......    
    // void PackBaseMsg(const std::string & sBodyBuffer, const int iCmd, std::string & sBuffer);
    // 修改代码
    int PackBaseMsg(const std::string & sBodyBuffer, const int iCmd, std::string & sBuffer);
    // ......  
};

src/algorithm/base.cpp

int Base :: PackMsg(const PaxosMsg & oPaxosMsg, std::string & sBuffer)
{
    std::string sBodyBuffer;
    bool bSucc = oPaxosMsg.SerializeToString(&sBodyBuffer);
    if (!bSucc)
    {   
        PLGErr("PaxosMsg.SerializeToString fail, skip this msg");
        return -1;
    }
    
    int iCmd = MsgCmd_PaxosMsg;
    // PackBaseMsg(sBodyBuffer, iCmd, sBuffer);
   
    // 修改代码
    int ret = PackBaseMsg(sBodyBuffer, iCmd, sBuffer);
    if (ret) {
        PLGErr("PaxosMsg PackBaseMsg fail");
        return ret;
    }
    
    return 0;
}

int Base :: PackCheckpointMsg(const CheckpointMsg & oCheckpointMsg, std::string & sBuffer)
{
    std::string sBodyBuffer;
    bool bSucc = oCheckpointMsg.SerializeToString(&sBodyBuffer);
    if (!bSucc)
    {
        PLGErr("CheckpointMsg.SerializeToString fail, skip this msg");
        return -1;
    }

    int iCmd = MsgCmd_CheckpointMsg;
    // PackBaseMsg(sBodyBuffer, iCmd, sBuffer);

    // 修改代码
    int ret = PackBaseMsg(sBodyBuffer, iCmd, sBuffer);
    if (ret) {
        PLGErr("CheckpointMsg PackBaseMsg fail");
        return ret;
    }     

    return 0;
}

// void Base :: PackBaseMsg(const std::string & sBodyBuffer, const int iCmd, std::string & sBuffer)
// 修改代码
int Base :: PackBaseMsg(const std::string & sBodyBuffer, const int iCmd, std::string & sBuffer)
{
    char sGroupIdx[GROUPIDXLEN] = {0};
    int iGroupIdx = m_poConfig->GetMyGroupIdx();
    memcpy(sGroupIdx, &iGroupIdx, sizeof(sGroupIdx));

    Header oHeader;
    oHeader.set_gid(m_poConfig->GetGid());
    oHeader.set_rid(0);
    oHeader.set_cmdid(iCmd);
    oHeader.set_version(1);

    std::string sHeaderBuffer;
    bool bSucc = oHeader.SerializeToString(&sHeaderBuffer);
    if (!bSucc)
    {
        PLGErr("Header.SerializeToString fail, skip this msg");
        // assert(bSucc == true);

        // 修改代码
        return -1;
    }

    char sHeaderLen[HEADLEN_LEN] = {0};
    uint16_t iHeaderLen = (uint16_t)sHeaderBuffer.size();
    memcpy(sHeaderLen, &iHeaderLen, sizeof(sHeaderLen));

    sBuffer = string(sGroupIdx, sizeof(sGroupIdx)) + string(sHeaderLen, sizeof(sHeaderLen)) + sHeaderBuffer + sBodyBuffer;

    //check sum
    uint32_t iBufferChecksum = crc32(0, (const uint8_t *)sBuffer.data(), sBuffer.size(), NET_CRC32SKIP);
    char sBufferChecksum[CHECKSUM_LEN] = {0};
    memcpy(sBufferChecksum, &iBufferChecksum, sizeof(sBufferChecksum));

    sBuffer += string(sBufferChecksum, sizeof(sBufferChecksum));

    // 新增代码
    return 0;
}

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions