| 1547 | } |
| 1548 | |
| 1549 | BlockingReason HashBuild::isBlocked(ContinueFuture* future) { |
| 1550 | switch (state_) { |
| 1551 | case State::kRunning: |
| 1552 | if (isInputFromSpill()) { |
| 1553 | processSpillInput(); |
| 1554 | } |
| 1555 | break; |
| 1556 | case State::kYield: |
| 1557 | setRunning(); |
| 1558 | BOLT_CHECK(isInputFromSpill()); |
| 1559 | processSpillInput(); |
| 1560 | break; |
| 1561 | case State::kFinish: |
| 1562 | break; |
| 1563 | case State::kWaitForSpill: |
| 1564 | if (!future_.valid()) { |
| 1565 | setRunning(); |
| 1566 | BOLT_CHECK_NOT_NULL(input_); |
| 1567 | DeltaCpuWallTimer timer{[this](const CpuWallTiming& timing) { |
| 1568 | auto selfDelta = |
| 1569 | operatorCtx_->driver()->processLazyTiming(*this, timing); |
| 1570 | this->stats().wlock()->addInputTiming.add(selfDelta); |
| 1571 | }}; |
| 1572 | addInput(std::move(input_)); |
| 1573 | } |
| 1574 | break; |
| 1575 | case State::kWaitForBuild: |
| 1576 | [[fallthrough]]; |
| 1577 | case State::kWaitForProbe: |
| 1578 | if (!future_.valid()) { |
| 1579 | setRunning(); |
| 1580 | postHashBuildProcess(); |
| 1581 | } |
| 1582 | break; |
| 1583 | default: |
| 1584 | BOLT_UNREACHABLE("Unexpected state: {}", stateName(state_)); |
| 1585 | break; |
| 1586 | } |
| 1587 | if (future_.valid()) { |
| 1588 | BOLT_CHECK(!isRunning() && !isFinished()); |
| 1589 | *future = std::move(future_); |
| 1590 | } |
| 1591 | return fromStateToBlockingReason(state_); |
| 1592 | } |
| 1593 | |
| 1594 | bool HashBuild::isFinished() { |
| 1595 | return state_ == State::kFinish; |
nothing calls this directly
no test coverage detected