a simple implementation of the DAG execution engine.
| 146 | |
| 147 | // a simple implementation of the DAG execution engine. |
| 148 | void 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; |
nothing calls this directly
no test coverage detected