| 701 | REGISTER_TASKFUNC(EraseLogRangeTaskFunc); |
| 702 | |
| 703 | struct CopyLogRangeTaskFunc : TaskFuncBase { |
| 704 | static StringRef name; |
| 705 | static constexpr uint32_t version = 1; |
| 706 | |
| 707 | static struct { |
| 708 | static TaskParam<int64_t> bytesWritten() { return LiteralStringRef(__FUNCTION__); } |
| 709 | } Params; |
| 710 | |
| 711 | static const Key keyNextBeginVersion; |
| 712 | |
| 713 | StringRef getName() const override { return name; }; |
| 714 | |
| 715 | Future<Void> execute(Database cx, |
| 716 | Reference<TaskBucket> tb, |
| 717 | Reference<FutureBucket> fb, |
| 718 | Reference<Task> task) override { |
| 719 | return _execute(cx, tb, fb, task); |
| 720 | }; |
| 721 | Future<Void> finish(Reference<ReadYourWritesTransaction> tr, |
| 722 | Reference<TaskBucket> tb, |
| 723 | Reference<FutureBucket> fb, |
| 724 | Reference<Task> task) override { |
| 725 | return _finish(tr, tb, fb, task); |
| 726 | }; |
| 727 | |
| 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()); |