| 159 | } |
| 160 | |
| 161 | void received(std::queue<Event> events) |
| 162 | { |
| 163 | while (!events.empty()) { |
| 164 | Event event = events.front(); |
| 165 | events.pop(); |
| 166 | |
| 167 | LOG(INFO) << "Received " << event.type() << " event"; |
| 168 | |
| 169 | switch (event.type()) { |
| 170 | case Event::SUBSCRIBED: { |
| 171 | framework.mutable_id()->CopyFrom(event.subscribed().framework_id()); |
| 172 | |
| 173 | LOG(INFO) << "Subscribed with ID '" << framework.id(); |
| 174 | state = SUBSCRIBED; |
| 175 | break; |
| 176 | } |
| 177 | |
| 178 | case Event::OFFERS: { |
| 179 | metrics.offers_received += event.offers().offers().size(); |
| 180 | |
| 181 | resourceOffers(google::protobuf::convert(event.offers().offers())); |
| 182 | break; |
| 183 | } |
| 184 | |
| 185 | case Event::INVERSE_OFFERS: { |
| 186 | metrics.inverse_offers_received += |
| 187 | event.inverse_offers().inverse_offers().size(); |
| 188 | |
| 189 | inverseOffers(google::protobuf::convert( |
| 190 | event.inverse_offers().inverse_offers())); |
| 191 | break; |
| 192 | } |
| 193 | |
| 194 | case Event::UPDATE: { |
| 195 | statusUpdate(event.update().status()); |
| 196 | break; |
| 197 | } |
| 198 | |
| 199 | // TODO(greggomann): Implement handling of operation status updates. |
| 200 | case Event::UPDATE_OPERATION_STATUS: |
| 201 | break; |
| 202 | |
| 203 | case Event::FAILURE: { |
| 204 | const Event::Failure& failure = event.failure(); |
| 205 | |
| 206 | if (failure.has_agent_id() && failure.has_executor_id()) { |
| 207 | LOG(INFO) |
| 208 | << "Executor '" << failure.executor_id() |
| 209 | << "' lost on agent '" << failure.agent_id() |
| 210 | << (failure.has_status() ? |
| 211 | "' with status: " + stringify(failure.status()) : ""); |
| 212 | } else { |
| 213 | CHECK(failure.has_agent_id()); |
| 214 | |
| 215 | LOG(INFO) << "Agent lost: " << failure.agent_id(); |
| 216 | } |
| 217 | break; |
| 218 | } |
nothing calls this directly
no test coverage detected