You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

527 lines
14 KiB

  1. package consensus
  2. import (
  3. "errors"
  4. "fmt"
  5. "sync"
  6. "time"
  7. cstypes "github.com/tendermint/tendermint/internal/consensus/types"
  8. tmsync "github.com/tendermint/tendermint/internal/libs/sync"
  9. "github.com/tendermint/tendermint/libs/bits"
  10. tmjson "github.com/tendermint/tendermint/libs/json"
  11. "github.com/tendermint/tendermint/libs/log"
  12. tmtime "github.com/tendermint/tendermint/libs/time"
  13. tmproto "github.com/tendermint/tendermint/proto/tendermint/types"
  14. "github.com/tendermint/tendermint/types"
  15. )
  16. var (
  17. ErrPeerStateHeightRegression = errors.New("peer state height regression")
  18. ErrPeerStateInvalidStartTime = errors.New("peer state invalid startTime")
  19. )
  20. // peerStateStats holds internal statistics for a peer.
  21. type peerStateStats struct {
  22. Votes int `json:"votes"`
  23. BlockParts int `json:"block_parts"`
  24. }
  25. func (pss peerStateStats) String() string {
  26. return fmt.Sprintf("peerStateStats{votes: %d, blockParts: %d}", pss.Votes, pss.BlockParts)
  27. }
  28. // PeerState contains the known state of a peer, including its connection and
  29. // threadsafe access to its PeerRoundState.
  30. // NOTE: THIS GETS DUMPED WITH rpc/core/consensus.go.
  31. // Be mindful of what you Expose.
  32. type PeerState struct {
  33. peerID types.NodeID
  34. logger log.Logger
  35. // NOTE: Modify below using setters, never directly.
  36. mtx tmsync.RWMutex
  37. running bool
  38. PRS cstypes.PeerRoundState `json:"round_state"`
  39. Stats *peerStateStats `json:"stats"`
  40. broadcastWG sync.WaitGroup
  41. closer *tmsync.Closer
  42. }
  43. // NewPeerState returns a new PeerState for the given node ID.
  44. func NewPeerState(logger log.Logger, peerID types.NodeID) *PeerState {
  45. return &PeerState{
  46. peerID: peerID,
  47. logger: logger,
  48. closer: tmsync.NewCloser(),
  49. PRS: cstypes.PeerRoundState{
  50. Round: -1,
  51. ProposalPOLRound: -1,
  52. LastCommitRound: -1,
  53. CatchupCommitRound: -1,
  54. },
  55. Stats: &peerStateStats{},
  56. }
  57. }
  58. // SetRunning sets the running state of the peer.
  59. func (ps *PeerState) SetRunning(v bool) {
  60. ps.mtx.Lock()
  61. defer ps.mtx.Unlock()
  62. ps.running = v
  63. }
  64. // IsRunning returns true if a PeerState is considered running where multiple
  65. // broadcasting goroutines exist for the peer.
  66. func (ps *PeerState) IsRunning() bool {
  67. ps.mtx.RLock()
  68. defer ps.mtx.RUnlock()
  69. return ps.running
  70. }
  71. // GetRoundState returns a shallow copy of the PeerRoundState. There's no point
  72. // in mutating it since it won't change PeerState.
  73. func (ps *PeerState) GetRoundState() *cstypes.PeerRoundState {
  74. ps.mtx.Lock()
  75. defer ps.mtx.Unlock()
  76. prs := ps.PRS.Copy()
  77. return &prs
  78. }
  79. // ToJSON returns a json of PeerState.
  80. func (ps *PeerState) ToJSON() ([]byte, error) {
  81. ps.mtx.Lock()
  82. defer ps.mtx.Unlock()
  83. return tmjson.Marshal(ps)
  84. }
  85. // GetHeight returns an atomic snapshot of the PeerRoundState's height used by
  86. // the mempool to ensure peers are caught up before broadcasting new txs.
  87. func (ps *PeerState) GetHeight() int64 {
  88. ps.mtx.Lock()
  89. defer ps.mtx.Unlock()
  90. return ps.PRS.Height
  91. }
  92. // SetHasProposal sets the given proposal as known for the peer.
  93. func (ps *PeerState) SetHasProposal(proposal *types.Proposal) {
  94. ps.mtx.Lock()
  95. defer ps.mtx.Unlock()
  96. if ps.PRS.Height != proposal.Height || ps.PRS.Round != proposal.Round {
  97. return
  98. }
  99. if ps.PRS.Proposal {
  100. return
  101. }
  102. ps.PRS.Proposal = true
  103. // ps.PRS.ProposalBlockParts is set due to NewValidBlockMessage
  104. if ps.PRS.ProposalBlockParts != nil {
  105. return
  106. }
  107. ps.PRS.ProposalBlockPartSetHeader = proposal.BlockID.PartSetHeader
  108. ps.PRS.ProposalBlockParts = bits.NewBitArray(int(proposal.BlockID.PartSetHeader.Total))
  109. ps.PRS.ProposalPOLRound = proposal.POLRound
  110. ps.PRS.ProposalPOL = nil // Nil until ProposalPOLMessage received.
  111. }
  112. // InitProposalBlockParts initializes the peer's proposal block parts header
  113. // and bit array.
  114. func (ps *PeerState) InitProposalBlockParts(partSetHeader types.PartSetHeader) {
  115. ps.mtx.Lock()
  116. defer ps.mtx.Unlock()
  117. if ps.PRS.ProposalBlockParts != nil {
  118. return
  119. }
  120. ps.PRS.ProposalBlockPartSetHeader = partSetHeader
  121. ps.PRS.ProposalBlockParts = bits.NewBitArray(int(partSetHeader.Total))
  122. }
  123. // SetHasProposalBlockPart sets the given block part index as known for the peer.
  124. func (ps *PeerState) SetHasProposalBlockPart(height int64, round int32, index int) {
  125. ps.mtx.Lock()
  126. defer ps.mtx.Unlock()
  127. if ps.PRS.Height != height || ps.PRS.Round != round {
  128. return
  129. }
  130. ps.PRS.ProposalBlockParts.SetIndex(index, true)
  131. }
  132. // PickVoteToSend picks a vote to send to the peer. It will return true if a
  133. // vote was picked.
  134. //
  135. // NOTE: `votes` must be the correct Size() for the Height().
  136. func (ps *PeerState) PickVoteToSend(votes types.VoteSetReader) (*types.Vote, bool) {
  137. ps.mtx.Lock()
  138. defer ps.mtx.Unlock()
  139. if votes.Size() == 0 {
  140. return nil, false
  141. }
  142. var (
  143. height = votes.GetHeight()
  144. round = votes.GetRound()
  145. votesType = tmproto.SignedMsgType(votes.Type())
  146. size = votes.Size()
  147. )
  148. // lazily set data using 'votes'
  149. if votes.IsCommit() {
  150. ps.ensureCatchupCommitRound(height, round, size)
  151. }
  152. ps.ensureVoteBitArrays(height, size)
  153. psVotes := ps.getVoteBitArray(height, round, votesType)
  154. if psVotes == nil {
  155. return nil, false // not something worth sending
  156. }
  157. if index, ok := votes.BitArray().Sub(psVotes).PickRandom(); ok {
  158. return votes.GetByIndex(int32(index)), true
  159. }
  160. return nil, false
  161. }
  162. func (ps *PeerState) getVoteBitArray(height int64, round int32, votesType tmproto.SignedMsgType) *bits.BitArray {
  163. if !types.IsVoteTypeValid(votesType) {
  164. return nil
  165. }
  166. if ps.PRS.Height == height {
  167. if ps.PRS.Round == round {
  168. switch votesType {
  169. case tmproto.PrevoteType:
  170. return ps.PRS.Prevotes
  171. case tmproto.PrecommitType:
  172. return ps.PRS.Precommits
  173. }
  174. }
  175. if ps.PRS.CatchupCommitRound == round {
  176. switch votesType {
  177. case tmproto.PrevoteType:
  178. return nil
  179. case tmproto.PrecommitType:
  180. return ps.PRS.CatchupCommit
  181. }
  182. }
  183. if ps.PRS.ProposalPOLRound == round {
  184. switch votesType {
  185. case tmproto.PrevoteType:
  186. return ps.PRS.ProposalPOL
  187. case tmproto.PrecommitType:
  188. return nil
  189. }
  190. }
  191. return nil
  192. }
  193. if ps.PRS.Height == height+1 {
  194. if ps.PRS.LastCommitRound == round {
  195. switch votesType {
  196. case tmproto.PrevoteType:
  197. return nil
  198. case tmproto.PrecommitType:
  199. return ps.PRS.LastCommit
  200. }
  201. }
  202. return nil
  203. }
  204. return nil
  205. }
  206. // 'round': A round for which we have a +2/3 commit.
  207. func (ps *PeerState) ensureCatchupCommitRound(height int64, round int32, numValidators int) {
  208. if ps.PRS.Height != height {
  209. return
  210. }
  211. /*
  212. NOTE: This is wrong, 'round' could change.
  213. e.g. if orig round is not the same as block LastCommit round.
  214. if ps.CatchupCommitRound != -1 && ps.CatchupCommitRound != round {
  215. panic(fmt.Sprintf(
  216. "Conflicting CatchupCommitRound. Height: %v,
  217. Orig: %v,
  218. New: %v",
  219. height,
  220. ps.CatchupCommitRound,
  221. round))
  222. }
  223. */
  224. if ps.PRS.CatchupCommitRound == round {
  225. return // Nothing to do!
  226. }
  227. ps.PRS.CatchupCommitRound = round
  228. if round == ps.PRS.Round {
  229. ps.PRS.CatchupCommit = ps.PRS.Precommits
  230. } else {
  231. ps.PRS.CatchupCommit = bits.NewBitArray(numValidators)
  232. }
  233. }
  234. // EnsureVoteBitArrays ensures the bit-arrays have been allocated for tracking
  235. // what votes this peer has received.
  236. // NOTE: It's important to make sure that numValidators actually matches
  237. // what the node sees as the number of validators for height.
  238. func (ps *PeerState) EnsureVoteBitArrays(height int64, numValidators int) {
  239. ps.mtx.Lock()
  240. defer ps.mtx.Unlock()
  241. ps.ensureVoteBitArrays(height, numValidators)
  242. }
  243. func (ps *PeerState) ensureVoteBitArrays(height int64, numValidators int) {
  244. if ps.PRS.Height == height {
  245. if ps.PRS.Prevotes == nil {
  246. ps.PRS.Prevotes = bits.NewBitArray(numValidators)
  247. }
  248. if ps.PRS.Precommits == nil {
  249. ps.PRS.Precommits = bits.NewBitArray(numValidators)
  250. }
  251. if ps.PRS.CatchupCommit == nil {
  252. ps.PRS.CatchupCommit = bits.NewBitArray(numValidators)
  253. }
  254. if ps.PRS.ProposalPOL == nil {
  255. ps.PRS.ProposalPOL = bits.NewBitArray(numValidators)
  256. }
  257. } else if ps.PRS.Height == height+1 {
  258. if ps.PRS.LastCommit == nil {
  259. ps.PRS.LastCommit = bits.NewBitArray(numValidators)
  260. }
  261. }
  262. }
  263. // RecordVote increments internal votes related statistics for this peer.
  264. // It returns the total number of added votes.
  265. func (ps *PeerState) RecordVote() int {
  266. ps.mtx.Lock()
  267. defer ps.mtx.Unlock()
  268. ps.Stats.Votes++
  269. return ps.Stats.Votes
  270. }
  271. // VotesSent returns the number of blocks for which peer has been sending us
  272. // votes.
  273. func (ps *PeerState) VotesSent() int {
  274. ps.mtx.Lock()
  275. defer ps.mtx.Unlock()
  276. return ps.Stats.Votes
  277. }
  278. // RecordBlockPart increments internal block part related statistics for this peer.
  279. // It returns the total number of added block parts.
  280. func (ps *PeerState) RecordBlockPart() int {
  281. ps.mtx.Lock()
  282. defer ps.mtx.Unlock()
  283. ps.Stats.BlockParts++
  284. return ps.Stats.BlockParts
  285. }
  286. // BlockPartsSent returns the number of useful block parts the peer has sent us.
  287. func (ps *PeerState) BlockPartsSent() int {
  288. ps.mtx.Lock()
  289. defer ps.mtx.Unlock()
  290. return ps.Stats.BlockParts
  291. }
  292. // SetHasVote sets the given vote as known by the peer
  293. func (ps *PeerState) SetHasVote(vote *types.Vote) {
  294. ps.mtx.Lock()
  295. defer ps.mtx.Unlock()
  296. ps.setHasVote(vote.Height, vote.Round, vote.Type, vote.ValidatorIndex)
  297. }
  298. func (ps *PeerState) setHasVote(height int64, round int32, voteType tmproto.SignedMsgType, index int32) {
  299. logger := ps.logger.With(
  300. "peerH/R", fmt.Sprintf("%d/%d", ps.PRS.Height, ps.PRS.Round),
  301. "H/R", fmt.Sprintf("%d/%d", height, round),
  302. )
  303. logger.Debug("setHasVote", "type", voteType, "index", index)
  304. // NOTE: some may be nil BitArrays -> no side effects
  305. psVotes := ps.getVoteBitArray(height, round, voteType)
  306. if psVotes != nil {
  307. psVotes.SetIndex(int(index), true)
  308. }
  309. }
  310. // ApplyNewRoundStepMessage updates the peer state for the new round.
  311. func (ps *PeerState) ApplyNewRoundStepMessage(msg *NewRoundStepMessage) {
  312. ps.mtx.Lock()
  313. defer ps.mtx.Unlock()
  314. // ignore duplicates or decreases
  315. if CompareHRS(msg.Height, msg.Round, msg.Step, ps.PRS.Height, ps.PRS.Round, ps.PRS.Step) <= 0 {
  316. return
  317. }
  318. var (
  319. psHeight = ps.PRS.Height
  320. psRound = ps.PRS.Round
  321. psCatchupCommitRound = ps.PRS.CatchupCommitRound
  322. psCatchupCommit = ps.PRS.CatchupCommit
  323. startTime = tmtime.Now().Add(-1 * time.Duration(msg.SecondsSinceStartTime) * time.Second)
  324. )
  325. ps.PRS.Height = msg.Height
  326. ps.PRS.Round = msg.Round
  327. ps.PRS.Step = msg.Step
  328. ps.PRS.StartTime = startTime
  329. if psHeight != msg.Height || psRound != msg.Round {
  330. ps.PRS.Proposal = false
  331. ps.PRS.ProposalBlockPartSetHeader = types.PartSetHeader{}
  332. ps.PRS.ProposalBlockParts = nil
  333. ps.PRS.ProposalPOLRound = -1
  334. ps.PRS.ProposalPOL = nil
  335. // we'll update the BitArray capacity later
  336. ps.PRS.Prevotes = nil
  337. ps.PRS.Precommits = nil
  338. }
  339. if psHeight == msg.Height && psRound != msg.Round && msg.Round == psCatchupCommitRound {
  340. // Peer caught up to CatchupCommitRound.
  341. // Preserve psCatchupCommit!
  342. // NOTE: We prefer to use prs.Precommits if
  343. // pr.Round matches pr.CatchupCommitRound.
  344. ps.PRS.Precommits = psCatchupCommit
  345. }
  346. if psHeight != msg.Height {
  347. // shift Precommits to LastCommit
  348. if psHeight+1 == msg.Height && psRound == msg.LastCommitRound {
  349. ps.PRS.LastCommitRound = msg.LastCommitRound
  350. ps.PRS.LastCommit = ps.PRS.Precommits
  351. } else {
  352. ps.PRS.LastCommitRound = msg.LastCommitRound
  353. ps.PRS.LastCommit = nil
  354. }
  355. // we'll update the BitArray capacity later
  356. ps.PRS.CatchupCommitRound = -1
  357. ps.PRS.CatchupCommit = nil
  358. }
  359. }
  360. // ApplyNewValidBlockMessage updates the peer state for the new valid block.
  361. func (ps *PeerState) ApplyNewValidBlockMessage(msg *NewValidBlockMessage) {
  362. ps.mtx.Lock()
  363. defer ps.mtx.Unlock()
  364. if ps.PRS.Height != msg.Height {
  365. return
  366. }
  367. if ps.PRS.Round != msg.Round && !msg.IsCommit {
  368. return
  369. }
  370. ps.PRS.ProposalBlockPartSetHeader = msg.BlockPartSetHeader
  371. ps.PRS.ProposalBlockParts = msg.BlockParts
  372. }
  373. // ApplyProposalPOLMessage updates the peer state for the new proposal POL.
  374. func (ps *PeerState) ApplyProposalPOLMessage(msg *ProposalPOLMessage) {
  375. ps.mtx.Lock()
  376. defer ps.mtx.Unlock()
  377. if ps.PRS.Height != msg.Height {
  378. return
  379. }
  380. if ps.PRS.ProposalPOLRound != msg.ProposalPOLRound {
  381. return
  382. }
  383. // TODO: Merge onto existing ps.PRS.ProposalPOL?
  384. // We might have sent some prevotes in the meantime.
  385. ps.PRS.ProposalPOL = msg.ProposalPOL
  386. }
  387. // ApplyHasVoteMessage updates the peer state for the new vote.
  388. func (ps *PeerState) ApplyHasVoteMessage(msg *HasVoteMessage) {
  389. ps.mtx.Lock()
  390. defer ps.mtx.Unlock()
  391. if ps.PRS.Height != msg.Height {
  392. return
  393. }
  394. ps.setHasVote(msg.Height, msg.Round, msg.Type, msg.Index)
  395. }
  396. // ApplyVoteSetBitsMessage updates the peer state for the bit-array of votes
  397. // it claims to have for the corresponding BlockID.
  398. // `ourVotes` is a BitArray of votes we have for msg.BlockID
  399. // NOTE: if ourVotes is nil (e.g. msg.Height < rs.Height),
  400. // we conservatively overwrite ps's votes w/ msg.Votes.
  401. func (ps *PeerState) ApplyVoteSetBitsMessage(msg *VoteSetBitsMessage, ourVotes *bits.BitArray) {
  402. ps.mtx.Lock()
  403. defer ps.mtx.Unlock()
  404. votes := ps.getVoteBitArray(msg.Height, msg.Round, msg.Type)
  405. if votes != nil {
  406. if ourVotes == nil {
  407. votes.Update(msg.Votes)
  408. } else {
  409. otherVotes := votes.Sub(ourVotes)
  410. hasVotes := otherVotes.Or(msg.Votes)
  411. votes.Update(hasVotes)
  412. }
  413. }
  414. }
  415. // String returns a string representation of the PeerState
  416. func (ps *PeerState) String() string {
  417. return ps.StringIndented("")
  418. }
  419. // StringIndented returns a string representation of the PeerState
  420. func (ps *PeerState) StringIndented(indent string) string {
  421. ps.mtx.Lock()
  422. defer ps.mtx.Unlock()
  423. return fmt.Sprintf(`PeerState{
  424. %s Key %v
  425. %s RoundState %v
  426. %s Stats %v
  427. %s}`,
  428. indent, ps.peerID,
  429. indent, ps.PRS.StringIndented(indent+" "),
  430. indent, ps.Stats,
  431. indent,
  432. )
  433. }