| 1652 | } |
| 1653 | |
| 1654 | void processIncomingMessages(std::vector<LogsDBRequest>& requests, std::vector<LogsDBResponse>& responses) { |
| 1655 | auto processingStarted = ternNow(); |
| 1656 | _maybeLogStatus(processingStarted); |
| 1657 | for(auto& resp : responses) { |
| 1658 | auto request = _reqResp.getRequest(resp.msg.id); |
| 1659 | if (request == nullptr) { |
| 1660 | // We often don't care about all responses and remove requests as soon as we can make progress |
| 1661 | continue; |
| 1662 | } |
| 1663 | |
| 1664 | // Mismatch in responses could be due to network issues we don't want to crash but we will ignore and retry |
| 1665 | // Mismatch in internal state is asserted on. |
| 1666 | if (unlikely(request->replicaId != resp.replicaId)) { |
| 1667 | LOG_ERROR(_env, "Expected response from replica %s, got it from replica %s. Response: %s", request->replicaId, resp.msg.id, resp); |
| 1668 | continue; |
| 1669 | } |
| 1670 | if (unlikely(request->msg.body.kind() != resp.msg.body.kind())) { |
| 1671 | LOG_ERROR(_env, "Expected response of type %s, got type %s. Response: %s", request->msg.body.kind(), resp.msg.body.kind(), resp); |
| 1672 | continue; |
| 1673 | } |
| 1674 | LOG_TRACE(_env, "processing %s", resp); |
| 1675 | |
| 1676 | switch(resp.msg.body.kind()) { |
| 1677 | case LogMessageKind::RELEASE: |
| 1678 | // We don't track release requests. This response is unexpected |
| 1679 | case LogMessageKind::ERROR: |
| 1680 | LOG_ERROR(_env, "Bad response %s", resp); |
| 1681 | break; |
| 1682 | case LogMessageKind::LOG_WRITE: |
| 1683 | _appender.proccessLogWriteResponse(request->replicaId, *request, resp.msg.body.getLogWrite()); |
| 1684 | break; |
| 1685 | case LogMessageKind::LOG_READ: |
| 1686 | _catchupReader.proccessLogReadResponse(request->replicaId, *request, resp.msg.body.getLogRead()); |
| 1687 | break; |
| 1688 | case LogMessageKind::NEW_LEADER: |
| 1689 | _leaderElection.proccessNewLeaderResponse(request->replicaId, *request, resp.msg.body.getNewLeader()); |
| 1690 | break; |
| 1691 | case LogMessageKind::NEW_LEADER_CONFIRM: |
| 1692 | _leaderElection.proccessNewLeaderConfirmResponse(request->replicaId, *request, resp.msg.body.getNewLeaderConfirm()); |
| 1693 | break; |
| 1694 | case LogMessageKind::LOG_RECOVERY_READ: |
| 1695 | _leaderElection.proccessRecoveryReadResponse(request->replicaId, *request, resp.msg.body.getLogRecoveryRead()); |
| 1696 | break; |
| 1697 | case LogMessageKind::LOG_RECOVERY_WRITE: |
| 1698 | _leaderElection.proccessRecoveryWriteResponse(request->replicaId, *request, resp.msg.body.getLogRecoveryWrite()); |
| 1699 | break; |
| 1700 | case LogMessageKind::EMPTY: |
| 1701 | ALWAYS_ASSERT("LogMessageKind::EMPTY should not happen"); |
| 1702 | break; |
| 1703 | } |
| 1704 | } |
| 1705 | for(auto& req : requests) { |
| 1706 | switch (req.msg.body.kind()) { |
| 1707 | case LogMessageKind::ERROR: |
| 1708 | LOG_ERROR(_env, "Bad request %s", req); |
| 1709 | break; |
| 1710 | case LogMessageKind::LOG_WRITE: |
| 1711 | _batchWriter.proccessLogWriteRequest(req); |
nothing calls this directly
no test coverage detected