| 25 | } |
| 26 | |
| 27 | void process(const MessagePtr& msg) { |
| 28 | NodeID sender = msg->sender; |
| 29 | Progress prog; |
| 30 | CHECK(prog.ParseFromString(msg->task.msg())); |
| 31 | if (merger_) { |
| 32 | merger_(prog, &progress_[sender]); |
| 33 | } else { |
| 34 | progress_[sender] = prog; |
| 35 | } |
| 36 | |
| 37 | double time = timer_.stop(); |
| 38 | if (time > interval_ && printer_) { |
| 39 | total_time_ += time; |
| 40 | printer_(total_time_, &progress_); |
| 41 | timer_.restart(); |
| 42 | } else { |
| 43 | timer_.start(); |
| 44 | } |
| 45 | } |
| 46 | private: |
| 47 | AllProgress progress_; |
| 48 | double interval_; |