276 lines
6.1 KiB
Go
276 lines
6.1 KiB
Go
package p2p
|
|
|
|
import (
|
|
"blockchain/chain"
|
|
"blockchain/config"
|
|
"fmt"
|
|
"math/rand"
|
|
"sync"
|
|
"time"
|
|
)
|
|
|
|
// link represents a network connection between two peers.
|
|
type link struct {
|
|
send chan Msg
|
|
latency time.Duration
|
|
}
|
|
|
|
func newlink(target Peer, latency time.Duration) *link {
|
|
l := &link{
|
|
send: make(chan Msg, 50),
|
|
latency: latency,
|
|
}
|
|
|
|
if l.latency > 0 {
|
|
panic("newlink: latency < 0")
|
|
}
|
|
|
|
go func() {
|
|
for msg := range l.send {
|
|
jitter := time.Duration(rand.NormFloat64()*float64(l.latency)/4 + float64(l.latency))
|
|
if jitter > 0 {
|
|
time.Sleep(jitter)
|
|
}
|
|
target.Receive(msg)
|
|
}
|
|
}()
|
|
|
|
return l
|
|
}
|
|
|
|
type Node struct {
|
|
Wallet *chain.Wallet
|
|
Chain *chain.Blockchain
|
|
Mempool *chain.Mempool
|
|
HashDelay time.Duration
|
|
conn chan Msg
|
|
syncCh chan Msg
|
|
peers map[Peer]*link
|
|
mx sync.Mutex
|
|
miner *miner
|
|
seen map[string]struct{}
|
|
}
|
|
|
|
func NewNode() *Node {
|
|
n := &Node{
|
|
Wallet: chain.NewWallet(),
|
|
Chain: chain.NewBlockchain(),
|
|
Mempool: chain.NewMempool(),
|
|
conn: make(chan Msg, 100),
|
|
syncCh: make(chan Msg, 10),
|
|
peers: make(map[Peer]*link),
|
|
seen: make(map[string]struct{}),
|
|
}
|
|
n.miner = newMiner(&n.HashDelay, func(block *chain.Block, took time.Duration) {
|
|
if n.Chain.AddBlock(block) {
|
|
n.Mempool.RemoveMined(block)
|
|
fmt.Printf("[%s] Mined block #%d (%d txs), hash: %s took=%v nonce=%d\n",
|
|
n.ShortID(), block.Height, len(block.Transactions), block.Hash[:16],
|
|
took.Round(time.Millisecond), block.Nonce)
|
|
emit(Event{Type: "block_mine", Node: n.ShortID(), BlockHash: block.Hash[:8], Height: block.Height, TxCount: len(block.Transactions)})
|
|
n.Broadcast(BlockMsg{Block: block, From: n})
|
|
} else {
|
|
fmt.Printf("[%s] MINE-REJECT #%d (%d txs) took=%v — mined valid PoW but AddBlock failed\n",
|
|
n.ShortID(), block.Height, len(block.Transactions), took.Round(time.Millisecond))
|
|
}
|
|
n.submitWork()
|
|
})
|
|
return n
|
|
}
|
|
|
|
// Receive implements Peer — delivers a network message.
|
|
func (n *Node) Receive(msg Msg) {
|
|
select {
|
|
case n.conn <- msg:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// ReceiveSync implements Peer — delivers a sync message.
|
|
func (n *Node) ReceiveSync(msg Msg) {
|
|
select {
|
|
case n.syncCh <- msg:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (n *Node) submitWork() {
|
|
last := n.Chain.LastBlock()
|
|
txs := n.Mempool.CollectBlock(n.Wallet.PublicKey)
|
|
block := chain.NewBlock(last.Height+1, txs, last.Hash)
|
|
n.miner.SubmitWork(block)
|
|
}
|
|
|
|
func (n *Node) ShortID() string {
|
|
return n.Wallet.ShortAddress()
|
|
}
|
|
|
|
func (n *Node) MarkSeen(hash string) bool {
|
|
n.mx.Lock()
|
|
defer n.mx.Unlock()
|
|
|
|
if _, ok := n.seen[hash]; ok {
|
|
return true
|
|
}
|
|
|
|
const maxSeen = 10000
|
|
|
|
if len(n.seen) >= maxSeen {
|
|
n.seen = make(map[string]struct{})
|
|
}
|
|
|
|
n.seen[hash] = struct{}{}
|
|
|
|
return false
|
|
}
|
|
|
|
func (n *Node) AddPeer(peer Peer) {
|
|
if peer == n {
|
|
return
|
|
}
|
|
n.mx.Lock()
|
|
defer n.mx.Unlock()
|
|
if len(n.peers) >= config.Default.P2P.MaxPeers {
|
|
return
|
|
}
|
|
if _, exists := n.peers[peer]; exists {
|
|
return
|
|
}
|
|
n.peers[peer] = newlink(peer, config.Default.P2P.NetDelay)
|
|
}
|
|
|
|
func (n *Node) PeerIDs() []string {
|
|
n.mx.Lock()
|
|
defer n.mx.Unlock()
|
|
ids := make([]string, 0, len(n.peers))
|
|
for p := range n.peers {
|
|
ids = append(ids, p.ShortID())
|
|
}
|
|
return ids
|
|
}
|
|
|
|
func (n *Node) snapshotPeers() []Peer {
|
|
n.mx.Lock()
|
|
defer n.mx.Unlock()
|
|
peers := make([]Peer, 0, len(n.peers))
|
|
for p := range n.peers {
|
|
peers = append(peers, p)
|
|
}
|
|
return peers
|
|
}
|
|
|
|
func (n *Node) Send(to Peer, msg Msg) {
|
|
n.mx.Lock()
|
|
link, ok := n.peers[to]
|
|
n.mx.Unlock()
|
|
if !ok {
|
|
return
|
|
}
|
|
select {
|
|
case link.send <- msg:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (n *Node) Broadcast(msg Msg) {
|
|
n.mx.Lock()
|
|
links := make([]*link, 0, len(n.peers))
|
|
for _, l := range n.peers {
|
|
links = append(links, l)
|
|
}
|
|
n.mx.Unlock()
|
|
|
|
for _, l := range links {
|
|
select {
|
|
case l.send <- msg:
|
|
default:
|
|
}
|
|
}
|
|
}
|
|
|
|
func (n *Node) Run(stop chan struct{}) {
|
|
go n.miner.Run(stop)
|
|
n.submitWork()
|
|
|
|
discover := time.NewTicker(10 * time.Second)
|
|
defer discover.Stop()
|
|
for {
|
|
select {
|
|
case <-discover.C:
|
|
peers := n.snapshotPeers()
|
|
if len(peers) > 0 {
|
|
n.Send(peers[rand.Intn(len(peers))], GetPeersMsg{From: n})
|
|
}
|
|
case msg := <-n.conn:
|
|
switch m := msg.(type) {
|
|
case BlockMsg:
|
|
if n.MarkSeen(m.Block.Hash) {
|
|
continue
|
|
}
|
|
if n.Chain.AddBlock(m.Block) {
|
|
n.Mempool.RemoveMined(m.Block)
|
|
fromID := ""
|
|
if m.From != nil {
|
|
fromID = m.From.ShortID()
|
|
}
|
|
emit(Event{Type: "block_accept", Node: n.ShortID(), BlockHash: m.Block.Hash[:8], Height: m.Block.Height, From: fromID, TxCount: len(m.Block.Transactions)})
|
|
n.Broadcast(BlockMsg{Block: m.Block, From: n})
|
|
n.submitWork()
|
|
} else if m.Block.Height > n.Chain.LastBlock().Height {
|
|
fromID := ""
|
|
if m.From != nil {
|
|
fromID = m.From.ShortID()
|
|
}
|
|
emit(Event{Type: "block_reject", Node: n.ShortID(), BlockHash: m.Block.Hash[:8], Height: m.Block.Height, From: fromID})
|
|
emit(Event{Type: "sync_request", Node: n.ShortID(), Height: n.Chain.LastBlock().Height, From: fromID})
|
|
n.requestSync(m.From)
|
|
}
|
|
case TxMsg:
|
|
if n.MarkSeen(m.Tx.TxID()) {
|
|
continue
|
|
}
|
|
n.Mempool.Add(m.Tx)
|
|
n.Broadcast(TxMsg{Tx: m.Tx, From: n})
|
|
case GetPeersMsg:
|
|
if m.From != nil {
|
|
n.Send(m.From, PeersMsg{Peers: n.snapshotPeers()})
|
|
}
|
|
case PeersMsg:
|
|
for _, p := range m.Peers {
|
|
n.AddPeer(p)
|
|
}
|
|
}
|
|
case msg := <-n.syncCh:
|
|
switch m := msg.(type) {
|
|
case GetBlocksMsg:
|
|
blocks := n.Chain.BlocksAfterHeight(m.Height)
|
|
if len(blocks) > 0 && m.From != nil {
|
|
m.From.ReceiveSync(BlocksMsg{Blocks: blocks})
|
|
}
|
|
case BlocksMsg:
|
|
oldHeight := n.Chain.LastBlock().Height
|
|
if n.Chain.ReplaceBlocks(m.Blocks) {
|
|
for _, b := range m.Blocks {
|
|
n.Mempool.RemoveMined(b)
|
|
}
|
|
newHeight := n.Chain.LastBlock().Height
|
|
emit(Event{Type: "reorg", Node: n.ShortID(), Height: newHeight, TxCount: newHeight - oldHeight})
|
|
fmt.Printf("[%s] Synced chain to block #%d\n",
|
|
n.ShortID(), newHeight)
|
|
n.submitWork()
|
|
}
|
|
}
|
|
case <-stop:
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
func (n *Node) requestSync(peer Peer) {
|
|
if peer == nil {
|
|
return
|
|
}
|
|
height := max(0, n.Chain.LastBlock().Height-5)
|
|
peer.ReceiveSync(GetBlocksMsg{Height: height, From: n})
|
|
}
|