| 424 | } |
| 425 | |
| 426 | func (db *DB) doWrites() { |
| 427 | defer db.statusMgr.Done() |
| 428 | pendingCh := make(chan struct{}, 1) |
| 429 | writeRequests := func(reqs []*request) { |
| 430 | if err := db.writeRequests(reqs); err != nil { |
| 431 | utils.GetLogger().Printf("writeRequests fail, err=%v", err) |
| 432 | panic(err) |
| 433 | } |
| 434 | <-pendingCh |
| 435 | } |
| 436 | |
| 437 | reqs := make([]*request, 0, 10) |
| 438 | var r *request |
| 439 | var ok bool |
| 440 | ctx := db.statusMgr.Context() |
| 441 | for { |
| 442 | select { |
| 443 | case <-ctx.Done(): |
| 444 | goto closedCase |
| 445 | case r, ok = <-db.writeCh: |
| 446 | if !ok { |
| 447 | goto closedCase |
| 448 | } |
| 449 | } |
| 450 | |
| 451 | for { |
| 452 | reqs = append(reqs, r) |
| 453 | |
| 454 | if len(reqs) >= 3*KvWriteChCapacity { |
| 455 | pendingCh <- struct{}{} // blocking. |
| 456 | goto writeCase |
| 457 | } |
| 458 | |
| 459 | select { |
| 460 | // Either push to pending, or continue to pick from writeCh. |
| 461 | case <-ctx.Done(): |
| 462 | goto closedCase |
| 463 | case r, ok = <-db.writeCh: |
| 464 | if !ok { |
| 465 | goto closedCase |
| 466 | } |
| 467 | case pendingCh <- struct{}{}: |
| 468 | goto writeCase |
| 469 | } |
| 470 | } |
| 471 | |
| 472 | closedCase: |
| 473 | // Drain pending requests and fail them, since we're shutting down. |
| 474 | for { |
| 475 | select { |
| 476 | case r = <-db.writeCh: |
| 477 | reqs = append(reqs, r) |
| 478 | default: |
| 479 | for _, req := range reqs { |
| 480 | req.Err = ErrDBClosed |
| 481 | req.Wg.Done() |
| 482 | } |
| 483 | return |