|
- package network
- import (
- "errors"
- "fmt"
- "net"
- "strconv"
- "time"
- "github.com/ethereum/go-ethereum/contracts/chequebook"
- "github.com/ethereum/go-ethereum/log"
- "github.com/ethereum/go-ethereum/metrics"
- "github.com/ethereum/go-ethereum/p2p"
- bzzswap "github.com/ethereum/go-ethereum/swarm/services/swap"
- "github.com/ethereum/go-ethereum/swarm/services/swap/swap"
- "github.com/ethereum/go-ethereum/swarm/storage"
- )
- var (
- storeRequestMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.storerequest.count", nil)
- retrieveRequestMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.retrieverequest.count", nil)
- peersMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.peers.count", nil)
- syncRequestMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.syncrequest.count", nil)
- unsyncedKeysMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.unsyncedkeys.count", nil)
- deliverRequestMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.deliverrequest.count", nil)
- paymentMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.payment.count", nil)
- invalidMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.invalid.count", nil)
- handleStatusMsgCounter = metrics.NewRegisteredCounter("network.protocol.msg.handlestatus.count", nil)
- )
- const (
- Version = 0
- ProtocolLength = uint64(8)
- ProtocolMaxMsgSize = 10 * 1024 * 1024
- NetworkId = 3
- )
- type bzz struct {
- storage StorageHandler
- hive *Hive
- dbAccess *DbAccess
- requestDb *storage.LDBDatabase
- remoteAddr *peerAddr
- peer *p2p.Peer
- rw p2p.MsgReadWriter
- backend chequebook.Backend
- lastActive time.Time
- NetworkId uint64
- swap *swap.Swap
- swapParams *bzzswap.SwapParams
- swapEnabled bool
- syncEnabled bool
- syncer *syncer
- syncParams *SyncParams
- syncState *syncState
- }
- type StorageHandler interface {
- HandleUnsyncedKeysMsg(req *unsyncedKeysMsgData, p *peer) error
- HandleDeliveryRequestMsg(req *deliveryRequestMsgData, p *peer) error
- HandleStoreRequestMsg(req *storeRequestMsgData, p *peer)
- HandleRetrieveRequestMsg(req *retrieveRequestMsgData, p *peer)
- }
- func Bzz(cloud StorageHandler, backend chequebook.Backend, hive *Hive, dbaccess *DbAccess, sp *bzzswap.SwapParams, sy *SyncParams, networkId uint64) (p2p.Protocol, error) {
-
-
- requestDb, err := storage.NewLDBDatabase(sy.RequestDbPath)
- if err != nil {
- return p2p.Protocol{}, fmt.Errorf("error setting up request db: %v", err)
- }
- if networkId == 0 {
- networkId = NetworkId
- }
- return p2p.Protocol{
- Name: "bzz",
- Version: Version,
- Length: ProtocolLength,
- Run: func(p *p2p.Peer, rw p2p.MsgReadWriter) error {
- return run(requestDb, cloud, backend, hive, dbaccess, sp, sy, networkId, p, rw)
- },
- }, nil
- }
- func run(requestDb *storage.LDBDatabase, depo StorageHandler, backend chequebook.Backend, hive *Hive, dbaccess *DbAccess, sp *bzzswap.SwapParams, sy *SyncParams, networkId uint64, p *p2p.Peer, rw p2p.MsgReadWriter) (err error) {
- self := &bzz{
- storage: depo,
- backend: backend,
- hive: hive,
- dbAccess: dbaccess,
- requestDb: requestDb,
- peer: p,
- rw: rw,
- swapParams: sp,
- syncParams: sy,
- swapEnabled: hive.swapEnabled,
- syncEnabled: true,
- NetworkId: networkId,
- }
-
- err = self.handleStatus()
- if err != nil {
- return err
- }
- defer func() {
-
-
- self.hive.removePeer(&peer{bzz: self})
- if self.syncer != nil {
- self.syncer.stop()
- }
- if self.swap != nil {
- self.swap.Stop()
- }
- }()
-
- for {
- if self.hive.blockRead {
- log.Warn(fmt.Sprintf("Cannot read network"))
- time.Sleep(100 * time.Millisecond)
- continue
- }
- err = self.handle()
- if err != nil {
- return
- }
- }
- }
- func (self *bzz) Drop() {
- self.peer.Disconnect(p2p.DiscSubprotocolError)
- }
- func (self *bzz) handle() error {
- msg, err := self.rw.ReadMsg()
- log.Debug(fmt.Sprintf("<- %v", msg))
- if err != nil {
- return err
- }
- if msg.Size > ProtocolMaxMsgSize {
- return fmt.Errorf("message too long: %v > %v", msg.Size, ProtocolMaxMsgSize)
- }
-
- defer msg.Discard()
- switch msg.Code {
- case statusMsg:
-
-
- log.Debug(fmt.Sprintf("Status message: %v", msg))
- return errors.New("extra status message")
- case storeRequestMsg:
-
- storeRequestMsgCounter.Inc(1)
- var req storeRequestMsgData
- if err := msg.Decode(&req); err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- if n := len(req.SData); n < 9 {
- return fmt.Errorf("<- %v: Data too short (%v)", msg, n)
- }
-
- self.lastActive = time.Now()
- log.Trace(fmt.Sprintf("incoming store request: %s", req.String()))
-
- self.storage.HandleStoreRequestMsg(&req, &peer{bzz: self})
- case retrieveRequestMsg:
-
- retrieveRequestMsgCounter.Inc(1)
- var req retrieveRequestMsgData
- if err := msg.Decode(&req); err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- req.from = &peer{bzz: self}
-
- if req.isLookup() {
- log.Trace(fmt.Sprintf("self lookup for %v: responding with peers only...", req.from))
- } else if req.Key == nil {
- return fmt.Errorf("protocol handler: req.Key == nil || req.Timeout == nil")
- } else {
-
- self.storage.HandleRetrieveRequestMsg(&req, &peer{bzz: self})
- }
-
- self.hive.peers(&req)
- case peersMsg:
-
-
- peersMsgCounter.Inc(1)
- var req peersMsgData
- if err := msg.Decode(&req); err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- req.from = &peer{bzz: self}
- log.Trace(fmt.Sprintf("<- peer addresses: %v", req))
- self.hive.HandlePeersMsg(&req, &peer{bzz: self})
- case syncRequestMsg:
- syncRequestMsgCounter.Inc(1)
- var req syncRequestMsgData
- if err := msg.Decode(&req); err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- log.Debug(fmt.Sprintf("<- sync request: %v", req))
- self.lastActive = time.Now()
- self.sync(req.SyncState)
- case unsyncedKeysMsg:
-
- unsyncedKeysMsgCounter.Inc(1)
- var req unsyncedKeysMsgData
- if err := msg.Decode(&req); err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- log.Debug(fmt.Sprintf("<- unsynced keys : %s", req.String()))
- err := self.storage.HandleUnsyncedKeysMsg(&req, &peer{bzz: self})
- self.lastActive = time.Now()
- if err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- case deliveryRequestMsg:
-
-
- deliverRequestMsgCounter.Inc(1)
- var req deliveryRequestMsgData
- if err := msg.Decode(&req); err != nil {
- return fmt.Errorf("<-msg %v: %v", msg, err)
- }
- log.Debug(fmt.Sprintf("<- delivery request: %s", req.String()))
- err := self.storage.HandleDeliveryRequestMsg(&req, &peer{bzz: self})
- self.lastActive = time.Now()
- if err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- case paymentMsg:
-
- paymentMsgCounter.Inc(1)
- if self.swapEnabled {
- var req paymentMsgData
- if err := msg.Decode(&req); err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- log.Debug(fmt.Sprintf("<- payment: %s", req.String()))
- self.swap.Receive(int(req.Units), req.Promise)
- }
- default:
-
- invalidMsgCounter.Inc(1)
- return fmt.Errorf("invalid message code: %v", msg.Code)
- }
- return nil
- }
- func (self *bzz) handleStatus() (err error) {
- handshake := &statusMsgData{
- Version: uint64(Version),
- ID: "honey",
- Addr: self.selfAddr(),
- NetworkId: self.NetworkId,
- Swap: &bzzswap.SwapProfile{
- Profile: self.swapParams.Profile,
- PayProfile: self.swapParams.PayProfile,
- },
- }
- err = p2p.Send(self.rw, statusMsg, handshake)
- if err != nil {
- return err
- }
-
- var msg p2p.Msg
- msg, err = self.rw.ReadMsg()
- if err != nil {
- return err
- }
- if msg.Code != statusMsg {
- return fmt.Errorf("first msg has code %x (!= %x)", msg.Code, statusMsg)
- }
- handleStatusMsgCounter.Inc(1)
- if msg.Size > ProtocolMaxMsgSize {
- return fmt.Errorf("message too long: %v > %v", msg.Size, ProtocolMaxMsgSize)
- }
- var status statusMsgData
- if err := msg.Decode(&status); err != nil {
- return fmt.Errorf("<- %v: %v", msg, err)
- }
- if status.NetworkId != self.NetworkId {
- return fmt.Errorf("network id mismatch: %d (!= %d)", status.NetworkId, self.NetworkId)
- }
- if Version != status.Version {
- return fmt.Errorf("protocol version mismatch: %d (!= %d)", status.Version, Version)
- }
- self.remoteAddr = self.peerAddr(status.Addr)
- log.Trace(fmt.Sprintf("self: advertised IP: %v, peer advertised: %v, local address: %v\npeer: advertised IP: %v, remote address: %v\n", self.selfAddr(), self.remoteAddr, self.peer.LocalAddr(), status.Addr.IP, self.peer.RemoteAddr()))
- if self.swapEnabled {
-
- self.swap, err = bzzswap.NewSwap(self.swapParams, status.Swap, self.backend, self)
- if err != nil {
- return err
- }
- }
- log.Info(fmt.Sprintf("Peer %08x is capable (%d/%d)", self.remoteAddr.Addr[:4], status.Version, status.NetworkId))
- err = self.hive.addPeer(&peer{bzz: self})
- if err != nil {
- return err
- }
-
- log.Info(fmt.Sprintf("syncronisation request sent with %v", self.syncState))
- self.syncRequest()
- return nil
- }
- func (self *bzz) sync(state *syncState) error {
-
- if self.syncer != nil {
- return errors.New("sync request can only be sent once")
- }
- cnt := self.dbAccess.counter()
- remoteaddr := self.remoteAddr.Addr
- start, stop := self.hive.kad.KeyRange(remoteaddr)
-
- if state == nil {
- self.syncEnabled = false
- log.Warn(fmt.Sprintf("syncronisation disabled for peer %v", self))
- state = &syncState{DbSyncState: &storage.DbSyncState{}, Synced: true}
- } else {
- state.synced = make(chan bool)
- state.SessionAt = cnt
- if storage.IsZeroKey(state.Stop) && state.Synced {
- state.Start = storage.Key(start[:])
- state.Stop = storage.Key(stop[:])
- }
- log.Debug(fmt.Sprintf("syncronisation requested by peer %v at state %v", self, state))
- }
- var err error
- self.syncer, err = newSyncer(
- self.requestDb,
- storage.Key(remoteaddr[:]),
- self.dbAccess,
- self.unsyncedKeys, self.store,
- self.syncParams, state, func() bool { return self.syncEnabled },
- )
- if err != nil {
- return nil
- }
- log.Trace(fmt.Sprintf("syncer set for peer %v", self))
- return nil
- }
- func (self *bzz) String() string {
- return self.remoteAddr.String()
- }
- func (self *bzz) peerAddr(base *peerAddr) *peerAddr {
- if base.IP.IsUnspecified() {
- host, _, _ := net.SplitHostPort(self.peer.RemoteAddr().String())
- base.IP = net.ParseIP(host)
- }
- return base
- }
- func (self *bzz) selfAddr() *peerAddr {
- id := self.hive.id
- host, port, _ := net.SplitHostPort(self.hive.listenAddr())
- intport, _ := strconv.Atoi(port)
- addr := &peerAddr{
- Addr: self.hive.addr,
- ID: id[:],
- IP: net.ParseIP(host),
- Port: uint16(intport),
- }
- return addr
- }
- func (self *bzz) retrieve(req *retrieveRequestMsgData) error {
- return self.send(retrieveRequestMsg, req)
- }
- func (self *bzz) store(req *storeRequestMsgData) error {
- return self.send(storeRequestMsg, req)
- }
- func (self *bzz) syncRequest() error {
- req := &syncRequestMsgData{}
- if self.hive.syncEnabled {
- log.Debug(fmt.Sprintf("syncronisation request to peer %v at state %v", self, self.syncState))
- req.SyncState = self.syncState
- }
- if self.syncState == nil {
- log.Warn(fmt.Sprintf("syncronisation disabled for peer %v at state %v", self, self.syncState))
- }
- return self.send(syncRequestMsg, req)
- }
- func (self *bzz) deliveryRequest(reqs []*syncRequest) error {
- req := &deliveryRequestMsgData{
- Deliver: reqs,
- }
- return self.send(deliveryRequestMsg, req)
- }
- func (self *bzz) unsyncedKeys(reqs []*syncRequest, state *syncState) error {
- req := &unsyncedKeysMsgData{
- Unsynced: reqs,
- State: state,
- }
- return self.send(unsyncedKeysMsg, req)
- }
- func (self *bzz) Pay(units int, promise swap.Promise) {
- req := &paymentMsgData{uint(units), promise.(*chequebook.Cheque)}
- self.payment(req)
- }
- func (self *bzz) payment(req *paymentMsgData) error {
- return self.send(paymentMsg, req)
- }
- func (self *bzz) peers(req *peersMsgData) error {
- return self.send(peersMsg, req)
- }
- func (self *bzz) send(msg uint64, data interface{}) error {
- if self.hive.blockWrite {
- return fmt.Errorf("network write blocked")
- }
- log.Trace(fmt.Sprintf("-> %v: %v (%T) to %v", msg, data, data, self))
- err := p2p.Send(self.rw, msg, data)
- if err != nil {
- self.Drop()
- }
- return err
- }
|