| 2068 | } |
| 2069 | |
| 2070 | void NameServerImpl::ProcessTask() { |
| 2071 | while (running_.load(std::memory_order_acquire)) { |
| 2072 | { |
| 2073 | bool has_task = false; |
| 2074 | std::unique_lock<std::mutex> lock(mu_); |
| 2075 | for (const auto& op_list : task_vec_) { |
| 2076 | if (!op_list.empty()) { |
| 2077 | has_task = true; |
| 2078 | break; |
| 2079 | } |
| 2080 | } |
| 2081 | if (!has_task) { |
| 2082 | cv_.wait_for(lock, std::chrono::milliseconds(FLAGS_name_server_task_wait_time)); |
| 2083 | if (!running_.load(std::memory_order_acquire)) { |
| 2084 | PDLOG(WARNING, "cur nameserver is not leader"); |
| 2085 | return; |
| 2086 | } |
| 2087 | } |
| 2088 | |
| 2089 | for (const auto& op_list : task_vec_) { |
| 2090 | if (op_list.empty()) { |
| 2091 | continue; |
| 2092 | } |
| 2093 | std::shared_ptr<OPData> op_data = op_list.front(); |
| 2094 | if (op_data->task_list_.empty() || op_data->op_info_.task_status() == ::fedb::api::kFailed || |
| 2095 | op_data->op_info_.task_status() == ::fedb::api::kCanceled) { |
| 2096 | continue; |
| 2097 | } |
| 2098 | if (op_data->op_info_.task_status() == ::fedb::api::kInited) { |
| 2099 | op_data->op_info_.set_start_time(::baidu::common::timer::now_time()); |
| 2100 | op_data->op_info_.set_task_status(::fedb::api::kDoing); |
| 2101 | std::string value; |
| 2102 | op_data->op_info_.SerializeToString(&value); |
| 2103 | std::string node = zk_op_data_path_ + "/" + std::to_string(op_data->op_info_.op_id()); |
| 2104 | if (!zk_client_->SetNodeValue(node, value)) { |
| 2105 | PDLOG(WARNING, "set zk op status value failed. node[%s] value[%s]", node.c_str(), |
| 2106 | value.c_str()); |
| 2107 | op_data->op_info_.set_task_status(::fedb::api::kInited); |
| 2108 | continue; |
| 2109 | } |
| 2110 | } |
| 2111 | std::shared_ptr<Task> task = op_data->task_list_.front(); |
| 2112 | if (task->task_info_->status() == ::fedb::api::kFailed) { |
| 2113 | PDLOG(WARNING, "task[%s] run failed, terminate op[%s]. op_id[%lu]", |
| 2114 | ::fedb::api::TaskType_Name(task->task_info_->task_type()).c_str(), |
| 2115 | ::fedb::api::OPType_Name(task->task_info_->op_type()).c_str(), task->task_info_->op_id()); |
| 2116 | } else if (task->task_info_->status() == ::fedb::api::kInited) { |
| 2117 | DEBUGLOG("run task. opid[%lu] op_type[%s] task_type[%s]", task->task_info_->op_id(), |
| 2118 | ::fedb::api::OPType_Name(task->task_info_->op_type()).c_str(), |
| 2119 | ::fedb::api::TaskType_Name(task->task_info_->task_type()).c_str()); |
| 2120 | task_thread_pool_.AddTask(task->fun_); |
| 2121 | task->task_info_->set_status(::fedb::api::kDoing); |
| 2122 | } else if (task->task_info_->status() == ::fedb::api::kDoing) { |
| 2123 | if (::baidu::common::timer::now_time() - op_data->op_info_.start_time() > |
| 2124 | FLAGS_name_server_op_execute_timeout / 1000) { |
| 2125 | PDLOG(INFO, |
| 2126 | "The execution time of op is too long. " |
| 2127 | "opid[%lu] op_type[%s] cur task_type[%s] " |
nothing calls this directly
no test coverage detected