(sender *blip.Sender, changeArray [][]interface{}, ignoreNoConflicts bool)
| 595 | } |
| 596 | |
| 597 | func (bh *blipHandler) sendBatchOfChanges(sender *blip.Sender, changeArray [][]interface{}, ignoreNoConflicts bool) error { |
| 598 | outrq := blip.NewRequest() |
| 599 | outrq.SetProfile("changes") |
| 600 | if ignoreNoConflicts { |
| 601 | outrq.Properties[ChangesMessageIgnoreNoConflicts] = trueProperty |
| 602 | } |
| 603 | if bh.collectionIdx != nil { |
| 604 | outrq.Properties[BlipCollection] = strconv.Itoa(*bh.collectionIdx) |
| 605 | } |
| 606 | err := outrq.SetJSONBody(changeArray) |
| 607 | if err != nil { |
| 608 | base.InfofCtx(bh.loggingCtx, base.KeyAll, "Error setting changes: %v", err) |
| 609 | } |
| 610 | |
| 611 | if len(changeArray) > 0 { |
| 612 | // Check for user updates before creating the db copy for handleChangesResponse |
| 613 | if err := bh.refreshUser(); err != nil { |
| 614 | return err |
| 615 | } |
| 616 | |
| 617 | handleChangesResponseDbCollection, err := bh.copyDatabaseCollectionWithUser(bh.collectionIdx) |
| 618 | if err != nil { |
| 619 | return err |
| 620 | } |
| 621 | |
| 622 | sendTime := time.Now() |
| 623 | if !bh.sendBLIPMessage(sender, outrq) { |
| 624 | return ErrClosedBLIPSender |
| 625 | } |
| 626 | |
| 627 | bh.inFlightChangesThrottle <- struct{}{} |
| 628 | atomic.AddInt64(&bh.changesPendingResponseCount, 1) |
| 629 | |
| 630 | bh.replicationStats.SendChangesCount.Add(int64(len(changeArray))) |
| 631 | // Spawn a goroutine to await the client's response: |
| 632 | go func(bh *blipHandler, sender *blip.Sender, response *blip.Message, changeArray [][]interface{}, sendTime time.Time, dbCollection *DatabaseCollectionWithUser) { |
| 633 | if err := bh.handleChangesResponse(bh.loggingCtx, sender, response, changeArray, sendTime, dbCollection, bh.collectionIdx); err != nil && !errors.Is(err, ErrClosedBLIPSender) { |
| 634 | base.WarnfCtx(bh.loggingCtx, "Error from bh.handleChangesResponse: %v", err) |
| 635 | if bh.fatalErrorCallback != nil { |
| 636 | bh.fatalErrorCallback(err) |
| 637 | } |
| 638 | } |
| 639 | |
| 640 | // Sent all of the revs for this changes batch, allow another changes batch to be sent. |
| 641 | select { |
| 642 | case <-bh.inFlightChangesThrottle: |
| 643 | case <-bh.terminator: |
| 644 | } |
| 645 | |
| 646 | atomic.AddInt64(&bh.changesPendingResponseCount, -1) |
| 647 | }(bh, sender, outrq.Response(), changeArray, sendTime, handleChangesResponseDbCollection) |
| 648 | } else { |
| 649 | outrq.SetNoReply(true) |
| 650 | if !bh.sendBLIPMessage(sender, outrq) { |
| 651 | return ErrClosedBLIPSender |
| 652 | } |
| 653 | } |
| 654 |
no test coverage detected