| 445 | } |
| 446 | |
| 447 | void ChunkServerImpl::WriteNextCallback(const WriteBlockRequest* next_request, |
| 448 | WriteBlockResponse* next_response, |
| 449 | bool failed, int error, |
| 450 | const std::string& next_server, |
| 451 | std::pair<const WriteBlockRequest*, WriteBlockResponse*> origin, |
| 452 | ::google::protobuf::Closure* done, |
| 453 | ChunkServer_Stub* stub) { |
| 454 | const WriteBlockRequest* request = origin.first; |
| 455 | WriteBlockResponse* response = origin.second; |
| 456 | /// If RPC_ERROR_SEND_BUFFER_FULL retry send. |
| 457 | if (failed && error == sofa::pbrpc::RPC_ERROR_SEND_BUFFER_FULL) { |
| 458 | std::function<void ()> callback = |
| 459 | std::bind(&ChunkServerImpl::WriteNext, this, next_server, |
| 460 | stub, next_request, next_response, request, response, done); |
| 461 | work_thread_pool_->DelayTask(10, callback); |
| 462 | return; |
| 463 | } |
| 464 | delete stub; |
| 465 | |
| 466 | int64_t block_id = request->block_id(); |
| 467 | const std::string& databuf = request->databuf(); |
| 468 | int64_t offset = request->offset(); |
| 469 | int32_t packet_seq = request->packet_seq(); |
| 470 | if (failed || next_response->status() != kOK) { |
| 471 | LOG(WARNING, "[WriteBlock] WriteNext %s fail: #%ld seq:%d, offset:%ld, len:%lu, " |
| 472 | "status= %s, error= %d\n", |
| 473 | next_server.c_str(), block_id, packet_seq, offset, databuf.size(), |
| 474 | StatusCode_Name(next_response->status()).c_str(), error); |
| 475 | if (failed) { |
| 476 | response->set_status(kNetworkUnavailable); |
| 477 | } else { |
| 478 | response->set_status(next_response->status()); |
| 479 | } |
| 480 | delete next_response; |
| 481 | g_unfinished_bytes.Sub(databuf.size()); |
| 482 | done->Run(); |
| 483 | return; |
| 484 | } else { |
| 485 | LOG(INFO, "[Writeblock] send #%ld seq:%d to next done", block_id, packet_seq); |
| 486 | delete next_response; |
| 487 | } |
| 488 | |
| 489 | std::function<void ()> callback = |
| 490 | std::bind(&ChunkServerImpl::LocalWriteBlock, this, request, response, done); |
| 491 | work_thread_pool_->AddTask(callback); |
| 492 | } |
| 493 | |
| 494 | void ChunkServerImpl::LocalWriteBlock(const WriteBlockRequest* request, |
| 495 | WriteBlockResponse* response, |