| 661 | } |
| 662 | } |
| 663 | void StartNextBlock() { |
| 664 | IUserContext::EDataDistrState ds = UserContext->UpdateDataDistrState(nullptr); |
| 665 | if (ds == IUserContext::DATA_UNAVAILABLE) { |
| 666 | Cancel(); |
| 667 | return; |
| 668 | } |
| 669 | int blockId, nextBlockId; |
| 670 | AtomicAdd(ActiveOpCount, 1); |
| 671 | if (ds == IUserContext::DATA_COMPLETE) { |
| 672 | int blockCount = Blocks.ysize(); |
| 673 | blockId = AtomicAdd(CurrentBlockId, blockCount) - blockCount; |
| 674 | nextBlockId = Blocks.ysize(); |
| 675 | } else { |
| 676 | blockId = AtomicAdd(CurrentBlockId, 1) - 1; |
| 677 | nextBlockId = blockId + 1; |
| 678 | } |
| 679 | if (blockId >= Blocks.ysize()) { |
| 680 | // no new ops to launch |
| 681 | // check if the work is complete |
| 682 | AtomicAdd(ActiveOpCount, -1); // cancel this block |
| 683 | if (AtomicAdd(ActiveOpCount, 0) == 0 && AtomicCas(&IsCanceledFlag, (void*)this, (void*)nullptr)) |
| 684 | TReduceExec::Launch(JobRequest.Get(), CompleteNotify.Get(), &ResultData, &ResultHasData); |
| 685 | } else { |
| 686 | int startIdx = Blocks[blockId].StartIdx; |
| 687 | int blockLen = 0; |
| 688 | for (int i = blockId; i < nextBlockId; ++i) |
| 689 | blockLen += Blocks[i].BlockLen; |
| 690 | LaunchBlockRequest(startIdx, blockLen); |
| 691 | } |
| 692 | } |
| 693 | void LaunchBlockRequest(int startIdx, int blockLen) { |
| 694 | TIntrusivePtr<TJobRequest> jr = new TJobRequest; |
| 695 | TVector<int> resultMap; |