| 566 | } |
| 567 | |
| 568 | void ProcessRpcRequest(InputMessageBase* msg_base) { |
| 569 | const int64_t start_parse_us = butil::cpuwide_time_us(); |
| 570 | DestroyingPtr<MostCommonMessage> msg(static_cast<MostCommonMessage*>(msg_base)); |
| 571 | SocketUniquePtr socket_guard(msg->ReleaseSocket()); |
| 572 | Socket* socket = socket_guard.get(); |
| 573 | const Server* server = static_cast<const Server*>(msg_base->arg()); |
| 574 | ScopedNonServiceError non_service_error(server); |
| 575 | |
| 576 | RpcMeta meta; |
| 577 | if (!ParsePbFromIOBuf(&meta, msg->meta)) { |
| 578 | LOG(WARNING) << "Fail to parse RpcMeta from " << *socket; |
| 579 | socket->SetFailed(EREQUEST, "Fail to parse RpcMeta from %s", |
| 580 | socket->description().c_str()); |
| 581 | return; |
| 582 | } |
| 583 | const RpcRequestMeta &request_meta = meta.request(); |
| 584 | |
| 585 | SampledRequest* sample = AskToBeSampled(); |
| 586 | if (sample) { |
| 587 | sample->meta.set_service_name(request_meta.service_name()); |
| 588 | sample->meta.set_method_name(request_meta.method_name()); |
| 589 | sample->meta.set_compress_type((CompressType)meta.compress_type()); |
| 590 | sample->meta.set_protocol_type(PROTOCOL_BAIDU_STD); |
| 591 | sample->meta.set_attachment_size(meta.attachment_size()); |
| 592 | sample->meta.set_authentication_data(meta.authentication_data()); |
| 593 | sample->request = msg->payload; |
| 594 | sample->submit(start_parse_us); |
| 595 | } |
| 596 | |
| 597 | std::unique_ptr<Controller> cntl(new (std::nothrow) Controller); |
| 598 | if (NULL == cntl.get()) { |
| 599 | LOG(WARNING) << "Fail to new Controller"; |
| 600 | return; |
| 601 | } |
| 602 | |
| 603 | RpcPBMessages* messages = NULL; |
| 604 | |
| 605 | ServerPrivateAccessor server_accessor(server); |
| 606 | ControllerPrivateAccessor accessor(cntl.get()); |
| 607 | const bool security_mode = server->options().security_mode() && |
| 608 | socket->user() == server_accessor.acceptor(); |
| 609 | if (request_meta.has_log_id()) { |
| 610 | cntl->set_log_id(request_meta.log_id()); |
| 611 | } |
| 612 | if (request_meta.has_request_id()) { |
| 613 | cntl->set_request_id(request_meta.request_id()); |
| 614 | } |
| 615 | if (request_meta.has_timeout_ms()) { |
| 616 | cntl->set_timeout_ms(request_meta.timeout_ms()); |
| 617 | } |
| 618 | cntl->set_request_content_type(meta.content_type()); |
| 619 | cntl->set_request_compress_type((CompressType)meta.compress_type()); |
| 620 | cntl->set_request_checksum_type((ChecksumType)meta.checksum_type()); |
| 621 | cntl->set_rpc_received_us(msg->received_us()); |
| 622 | accessor.set_checksum_value(meta.checksum_value()); |
| 623 | accessor.set_server(server) |
| 624 | .set_security_mode(security_mode) |
| 625 | .set_peer_id(socket->id()) |
no test coverage detected