@braft::StateMachine
| 251 | |
| 252 | // @braft::StateMachine |
| 253 | void on_apply(braft::Iterator& iter) { |
| 254 | // A batch of tasks are committed, which must be processed through |
| 255 | // |iter| |
| 256 | for (; iter.valid(); iter.next()) { |
| 257 | BlockResponse* response = NULL; |
| 258 | // This guard helps invoke iter.done()->Run() asynchronously to |
| 259 | // avoid that callback blocks the StateMachine |
| 260 | braft::AsyncClosureGuard closure_guard(iter.done()); |
| 261 | butil::IOBuf data; |
| 262 | off_t offset = 0; |
| 263 | if (iter.done()) { |
| 264 | // This task is applied by this node, get value from this |
| 265 | // closure to avoid additional parsing. |
| 266 | BlockClosure* c = dynamic_cast<BlockClosure*>(iter.done()); |
| 267 | offset = c->request()->offset(); |
| 268 | data.swap(*(c->data())); |
| 269 | response = c->response(); |
| 270 | } else { |
| 271 | // Have to parse BlockRequest from this log. |
| 272 | uint32_t meta_size = 0; |
| 273 | butil::IOBuf saved_log = iter.data(); |
| 274 | saved_log.cutn(&meta_size, sizeof(uint32_t)); |
| 275 | // Remember that meta_size is in network order which hould be |
| 276 | // covert to host order |
| 277 | meta_size = butil::NetToHost32(meta_size); |
| 278 | butil::IOBuf meta; |
| 279 | saved_log.cutn(&meta, meta_size); |
| 280 | butil::IOBufAsZeroCopyInputStream wrapper(meta); |
| 281 | BlockRequest request; |
| 282 | CHECK(request.ParseFromZeroCopyStream(&wrapper)); |
| 283 | data.swap(saved_log); |
| 284 | offset = request.offset(); |
| 285 | } |
| 286 | |
| 287 | const ssize_t nw = braft::file_pwrite(data, _fd->fd(), offset); |
| 288 | if (nw < 0) { |
| 289 | PLOG(ERROR) << "Fail to write to fd=" << _fd->fd(); |
| 290 | if (response) { |
| 291 | response->set_success(false); |
| 292 | } |
| 293 | // Let raft run this closure. |
| 294 | closure_guard.release(); |
| 295 | // Some disk error occurred, notify raft and never apply any data |
| 296 | // ever after |
| 297 | iter.set_error_and_rollback(); |
| 298 | return; |
| 299 | } |
| 300 | |
| 301 | if (response) { |
| 302 | response->set_success(true); |
| 303 | } |
| 304 | |
| 305 | // The purpose of following logs is to help you understand the way |
| 306 | // this StateMachine works. |
| 307 | // Remove these logs in performance-sensitive servers. |
| 308 | LOG_IF(INFO, FLAGS_log_applied_task) |
| 309 | << "Write " << data.size() << " bytes" |
| 310 | << " from offset=" << offset |
nothing calls this directly
no test coverage detected