| 15 | RpcSchedule::~RpcSchedule() { delete policy_; } |
| 16 | |
| 17 | void RpcSchedule::EnqueueRpc(const std::string& table_name, RpcTask* rpc) { |
| 18 | MutexLock lock(&mutex_); |
| 19 | |
| 20 | ScheduleEntity* entity = NULL; |
| 21 | TableList::iterator it = table_list_.find(table_name); |
| 22 | if (it != table_list_.end()) { |
| 23 | entity = it->second; |
| 24 | } else { |
| 25 | entity = table_list_[table_name] = policy_->NewScheEntity(new TaskQueue); |
| 26 | } |
| 27 | |
| 28 | TaskQueue* task_queue = (TaskQueue*)entity->user_ptr; |
| 29 | task_queue->push(rpc); |
| 30 | |
| 31 | task_queue->pending_count++; |
| 32 | pending_task_count_++; |
| 33 | |
| 34 | if (task_queue->pending_count == 1) { |
| 35 | policy_->Enable(entity); |
| 36 | } |
| 37 | } |
| 38 | |
| 39 | bool RpcSchedule::DequeueRpc(RpcTask** rpc) { |
| 40 | MutexLock lock(&mutex_); |
no test coverage detected