MCPcopy Create free account
hub / github.com/XTXMarkets/ternfs / processIncomingMessages

Method processIncomingMessages

cpp/core/LogsDB.cpp:1654–1750  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

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);

Callers

nothing calls this directly

Calls 15

ternNowFunction · 0.85
update_atomic_stat_emaFunction · 0.85
getRequestMethod · 0.80

Tested by

no test coverage detected