MCPcopy Create free account
hub / github.com/nutsdb/nutsdb / doWrites

Method doWrites

db.go:426–491  ·  view source on GitHub ↗
()

Source from the content-addressed store, hash-verified

424}
425
426func (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

Callers 1

openFunction · 0.95

Calls 6

writeRequestsMethod · 0.95
GetLoggerFunction · 0.92
appendFunction · 0.85
PrintfMethod · 0.65
DoneMethod · 0.45
ContextMethod · 0.45

Tested by

no test coverage detected