| 60 | } |
| 61 | |
| 62 | void Thread::PollMsgChannel() { |
| 63 | while (true) { |
| 64 | if (local_msg_queue_.empty()) { |
| 65 | CHECK_EQ(msg_channel_.ReceiveMany(&local_msg_queue_), kChannelStatusSuccess); |
| 66 | } |
| 67 | ActorMsg msg = std::move(local_msg_queue_.front()); |
| 68 | local_msg_queue_.pop(); |
| 69 | if (msg.msg_type() == ActorMsgType::kCmdMsg) { |
| 70 | if (msg.actor_cmd() == ActorCmd::kStopThread) { |
| 71 | CHECK(id2actor_ptr_.empty()) |
| 72 | << " RuntimeError! Thread: " << thrd_id_ |
| 73 | << " NOT empty when stop with actor num: " << id2actor_ptr_.size(); |
| 74 | break; |
| 75 | } else if (msg.actor_cmd() == ActorCmd::kConstructActor) { |
| 76 | ConstructActor(msg.dst_actor_id()); |
| 77 | continue; |
| 78 | } else { |
| 79 | // do nothing |
| 80 | } |
| 81 | } |
| 82 | int64_t actor_id = msg.dst_actor_id(); |
| 83 | auto actor_it = id2actor_ptr_.find(actor_id); |
| 84 | CHECK(actor_it != id2actor_ptr_.end()); |
| 85 | int process_msg_ret = actor_it->second.second->ProcessMsg(msg); |
| 86 | if (process_msg_ret == 1) { |
| 87 | VLOG(3) << "thread " << thrd_id_ << " deconstruct actor " << actor_id; |
| 88 | auto job_id_it = id2job_id_.find(actor_id); |
| 89 | const int64_t job_id = job_id_it->second; |
| 90 | id2job_id_.erase(job_id_it); |
| 91 | id2actor_ptr_.erase(actor_it); |
| 92 | Singleton<RuntimeCtx>::Get()->DecreaseCounter(GetRunningActorCountKeyByJobId(job_id)); |
| 93 | } else { |
| 94 | CHECK_EQ(process_msg_ret, 0); |
| 95 | } |
| 96 | } |
| 97 | } |
| 98 | |
| 99 | void Thread::ConstructActor(int64_t actor_id) { |
| 100 | std::unique_lock<std::mutex> lck(id2task_mtx_); |
nothing calls this directly
no test coverage detected