| 728 | // store mutation data from results until the end of stream or the timeout. If breaks on timeout returns the first |
| 729 | // uncopied version |
| 730 | ACTOR static Future<Optional<Version>> dumpData(Database cx, |
| 731 | Reference<Task> task, |
| 732 | PromiseStream<RCGroup> results, |
| 733 | FlowLock* lock, |
| 734 | Reference<TaskBucket> tb, |
| 735 | double breakTime) { |
| 736 | state bool endOfStream = false; |
| 737 | state Subspace conf = Subspace(databaseBackupPrefixRange.begin) |
| 738 | .get(BackupAgentBase::keyConfig) |
| 739 | .get(task->params[BackupAgentBase::keyConfigLogUid]); |
| 740 | state std::vector<RangeResult> nextMutations; |
| 741 | state bool isTimeoutOccured = false; |
| 742 | state Optional<KeyRef> lastKey; |
| 743 | state Version lastVersion; |
| 744 | state int64_t nextMutationSize = 0; |
| 745 | loop { |
| 746 | try { |
| 747 | if (endOfStream && !nextMutationSize) { |
| 748 | return Optional<Version>(); |
| 749 | } |
| 750 | |
| 751 | state std::vector<RangeResult> mutations = std::move(nextMutations); |
| 752 | state int64_t mutationSize = nextMutationSize; |
| 753 | nextMutations = std::vector<RangeResult>(); |
| 754 | nextMutationSize = 0; |
| 755 | |
| 756 | if (!endOfStream) { |
| 757 | loop { |
| 758 | try { |
| 759 | RCGroup group = waitNext(results.getFuture()); |
| 760 | lock->release(group.items.expectedSize()); |
| 761 | |
| 762 | int vecSize = group.items.expectedSize(); |
| 763 | if (mutationSize + vecSize >= CLIENT_KNOBS->BACKUP_LOG_WRITE_BATCH_MAX_SIZE) { |
| 764 | |
| 765 | nextMutations.push_back(group.items); |
| 766 | nextMutationSize = vecSize; |
| 767 | break; |
| 768 | } |
| 769 | |
| 770 | mutations.push_back(group.items); |
| 771 | mutationSize += vecSize; |
| 772 | } catch (Error& e) { |
| 773 | state Error error = e; |
| 774 | if (e.code() == error_code_end_of_stream) { |
| 775 | endOfStream = true; |
| 776 | break; |
| 777 | } |
| 778 | |
| 779 | throw error; |
| 780 | } |
| 781 | } |
| 782 | } |
| 783 | |
| 784 | state Optional<Version> nextVersionAfterBreak; |
| 785 | state Transaction tr(cx); |
| 786 | |
| 787 | loop { |
nothing calls this directly
no test coverage detected