MCPcopy Create free account
hub / github.com/Oneflow-Inc/oneflow / PollMsgChannel

Method PollMsgChannel

oneflow/core/thread/thread.cpp:62–97  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

60}
61
62void 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
99void Thread::ConstructActor(int64_t actor_id) {
100 std::unique_lock<std::mutex> lck(id2task_mtx_);

Callers

nothing calls this directly

Calls 13

GetFunction · 0.85
ReceiveManyMethod · 0.80
msg_typeMethod · 0.80
actor_cmdMethod · 0.80
dst_actor_idMethod · 0.80
findMethod · 0.80
DecreaseCounterMethod · 0.80
emptyMethod · 0.45
popMethod · 0.45
sizeMethod · 0.45
endMethod · 0.45

Tested by

no test coverage detected