MCPcopy Create free account
hub / github.com/dmlc/parameter_server / run

Method run

src/system/executor.cc:148–267  ·  view source on GitHub ↗

a simple implementation of the DAG execution engine.

Source from the content-addressed store, hash-verified

146
147// a simple implementation of the DAG execution engine.
148void Executor::run() {
149 while (!done_) {
150 bool process = false;
151 if (FLAGS_verbose) {
152 LI << obj_.name() << " before entering task loop; recved_msgs_. size [" <<
153 recved_msgs_.size() << "]";
154 }
155 {
156 std::unique_lock<std::mutex> lk(recved_msg_mu_);
157 // pickup a message with dependency satisfied
158 for (auto it = recved_msgs_.begin(); it != recved_msgs_.end(); ++it) {
159 auto& msg = *it;
160 auto sender = rnode(msg->sender);
161 if (!sender) {
162 LL << my_node_.id() << ": " << msg->sender
163 << " does not exist, ignore\n" << msg->debugString();
164 recved_msgs_.erase(it);
165 break;
166 }
167 // ack message, no dependency constraint
168 process = !msg->task.request();
169 if (!process) {
170 // check if the dependency constraints are satisfied
171 bool satisfied = true;
172 for (int i = 0; i < msg->task.wait_time_size(); ++i) {
173 int wait_time = msg->task.wait_time(i);
174 if (wait_time > Message::kInvalidTime &&
175 !sender->tryWaitIncomingTask(wait_time)) {
176 satisfied = false;
177 }
178 }
179 process = satisfied;
180 }
181 if (process) {
182 active_msg_ = msg;
183 recved_msgs_.erase(it);
184 if (FLAGS_verbose) {
185 LI << obj_.name() << " picked up an active_msg_ from recved_msgs_. " <<
186 "remaining size [" << recved_msgs_.size() << "]; msg [" <<
187 active_msg_->shortDebugString() << "]";
188 }
189 break;
190 }
191 }
192 if (!process) {
193 if (FLAGS_verbose) {
194 LI << obj_.name() << " picked nothing from recved_msgs_. size [" <<
195 recved_msgs_.size() << "] waiting Executor::accept";
196 }
197 dag_cond_.wait(lk);
198 continue;
199 }
200 }
201 // process the picked message
202 bool req = active_msg_->task.request();
203 int t = active_msg_->task.time();
204 auto sender = rnode(active_msg_->sender);
205 CHECK(sender) << "unknow node: " << active_msg_->sender;

Callers

nothing calls this directly

Calls 15

requestMethod · 0.80
wait_time_sizeMethod · 0.80
wait_timeMethod · 0.80
tryWaitIncomingTaskMethod · 0.80
waitMethod · 0.80
replyMethod · 0.80
tryWaitOutgoingTaskMethod · 0.80
sizeMethod · 0.45
beginMethod · 0.45
endMethod · 0.45
idMethod · 0.45
debugStringMethod · 0.45

Tested by

no test coverage detected